Stream Processing and Event Streaming Questions
Building on event-streaming platforms: Kafka and message queues, event sourcing, partitioning, consumer groups, exactly-once vs at-least-once delivery, and windowing. Covers handling late and out-of-order events, watermarks, and stateful stream operators. The core skill for real-time data engineering.
Design the operational controls for a multi-tenant streaming platform (shared Kafka cluster and stream-processing jobs) used by many teams: resource quotas, access controls, network isolation, and how you'd prevent one noisy tenant from degrading others.
Sample Answer
Direct answer
Operating a multi-tenant streaming platform means enforcing per-tenant resource quotas, isolating access via authentication and authorization, and building in guardrails so one tenant's misbehaving workload can't degrade everyone else's throughput or latency.
Structured elaboration
Resource quotas (throughput quotas on produce/fetch requests, storage quotas per topic or tenant) are the primary lever against noisy neighbors: without them, one tenant pushing an unexpected volume spike can consume shared broker bandwidth or CPU that every other tenant depends on. Access controls (ACLs scoping which principals can produce or consume from which topics, network-level isolation for genuinely sensitive tenants) prevent one team from accidentally (or maliciously) reading or writing another's data. For stateful stream-processing jobs specifically (Flink or Kafka Streams jobs run on shared compute, often on Kubernetes), the same isolation problem shows up as noisy-neighbor CPU/memory contention between jobs on the same cluster, addressed with resource requests/limits per job, dedicated node pools for the largest tenants, and safe rolling-upgrade paths that don't require taking the whole shared cluster down.
Worked example
Concretely, a shared cluster serving five teams would set a per-tenant produce-byte-rate quota (so no single tenant's traffic spike starves the others), ACLs scoping each team's service accounts to only their own topic prefix, and, for the Kubernetes-hosted stream-processing jobs layer, per-job resource requests/limits plus pod anti-affinity rules so two large jobs don't land on the same node and compete for the same CPU. A safe upgrade path (rolling one node at a time, verifying job health before proceeding) avoids a platform-wide outage when applying a change that affects every tenant's jobs.
Trade-offs and pitfalls
Too-strict quotas can throttle a legitimate tenant during an expected traffic spike (a marketing campaign, a seasonal peak) unless quotas are reviewed and adjusted proactively rather than reactively. The common operational mistake is treating multi-tenancy as purely an access-control problem and skipping resource quotas entirely, which leaves the platform vulnerable to exactly the noisy-neighbor scenario multi-tenancy is supposed to guard against; ACLs stop the wrong people from touching the wrong data, but they don't stop the right person's workload from starving everyone else.
How would you diagnose and respond to network partitions and split-brain-like symptoms in a streaming ecosystem where brokers, a controller/coordination layer, and downstream stream processors all disagree about cluster state?
Sample Answer
Direct answer
Network partitions and split-brain-like symptoms in a streaming ecosystem show up as brokers, their controller/coordination layer, and downstream processors disagreeing about who's the leader for a partition or which consumer owns which work; the fix is always to trust the coordination layer's quorum decision and force everything else (leader assignment, consumer offsets) to converge to it, never to let two sides keep operating independently.
Structured elaboration
A network partition can split a broker cluster into two groups that can each see a majority (or, worse, neither sees a majority) of the coordination service. If the coordination layer correctly demotes a minority-side broker (a broker that can no longer reach quorum steps down as leader for its partitions), the system self-heals once connectivity is restored: the minority side simply catches up as a follower. The dangerous case is a controller or coordination-layer bug (or misconfiguration like unclean leader election enabled) that lets a minority-side broker keep acting as leader, accepting writes that the majority side never sees, which produces genuinely divergent histories for the same partition, a real split-brain.
Worked example
Concretely, your triage sequence: first check the coordination layer's own health and quorum status (is it itself partitioned, or reporting a clean picture); then check for any partition whose reported leader differs between what different brokers believe versus what the coordination layer says, which is the direct signature of a stale or diverged leader; then check consumer groups for members that appear registered from two different broker views simultaneously, an equivalent split-brain at the consumer-group level. Once found, the resolution is to force the ecosystem to converge on the coordination layer's authoritative view: fence off (or restart) any broker or processor still operating on a stale view, and, if unclean leader election was involved, treat the affected partitions' recent history as suspect and reconcile against a known-good replica.
Trade-offs and pitfalls
Disabling unclean leader election trades availability for consistency: if the only in-sync replica is unreachable, the partition simply goes unavailable rather than electing a stale replica that could silently lose recently acknowledged writes. Teams under availability pressure sometimes enable unclean leader election broadly to avoid outages, without recognizing they've traded a temporary availability problem for a permanent, silent data-consistency problem that's much harder to detect after the fact.
What is tiered storage for a commit-log platform, and how does offloading older log segments to object storage change broker disk usage, read/write latency, and the economics of long-retention or replay-heavy topics?
Sample Answer
Direct answer
Tiered storage moves older log segments from local broker disk to cheaper object storage, which shrinks the broker's local storage footprint and cost dramatically at the price of higher latency for reads that reach back into that older, offloaded data.
Structured elaboration
Without tiered storage, a broker's local disk has to hold every byte of retained data for every partition it leads, which caps how long you can afford to retain data (or requires expensive, large local disks). Tiered storage splits the log into a "hot" tier still on fast local disk (recent segments, actively being written and read) and a "cold" tier offloaded to object storage (older, closed segments). Reads for recent data are unaffected; reads reaching into the cold tier incur the latency of fetching from object storage, generally acceptable for the kind of use case tiered storage targets: long-retention or replay-heavy topics, not latency-critical hot-path consumption.
Worked example
A topic needing 2 years of retention for regulatory replay, but where 99% of real read traffic only ever touches the last 24 hours, is the textbook fit: local disks only need to hold roughly a day's worth of data (dramatically reducing the broker fleet's total disk footprint and cost) while object storage cheaply holds the other 729 days, accessed only on the rare occasion someone needs to replay old history.
Trade-offs and pitfalls
A workload that frequently reads far back into history (contradicting the "cold data is rarely read" assumption tiered storage is built around) will see a real latency and cost hit from constant object-storage fetches, potentially worse than just provisioning enough local disk in the first place. Tiered storage also adds an operational dependency on the object-storage layer's availability for any read that reaches the cold tier, which is a new failure mode a purely local-disk deployment didn't have.
How does a distributed commit-log platform like Kafka differ from a traditional message broker such as RabbitMQ? Discuss replayability, multiple independent consumers reading the same data, and when you'd still prefer a traditional broker.
Sample Answer
Direct answer
A distributed commit-log platform stores every message durably and lets many independent consumer groups replay the same data at their own pace, while a traditional message broker typically removes a message once it's been acknowledged by a consumer, favoring simple point-to-point delivery over replay.
Structured elaboration
In a traditional broker, a message is usually a transient unit of work: once consumed and acknowledged, it's gone, and if you need a second, independent consumer to see the same messages, you generally need a separate queue or a fan-out mechanism the broker sets up. A commit-log platform instead retains messages for a configured window (or forever, with compaction) regardless of whether anyone has consumed them, and any number of independent consumer groups can each read the same log at their own offset and pace without affecting each other. This makes replay (reprocessing history, backfilling a new downstream system, recovering from a bug) natural in a commit-log system and comparatively awkward in a pure queue.
Worked example
Concretely: an event-sourced billing system wants both a real-time fraud-scoring consumer and, three months later, a brand-new analytics team that needs to reprocess the last 30 days of billing events to build a new report. On a commit-log platform, the new consumer group simply starts reading from an earlier offset (or from the retention window's start) with zero coordination with the existing fraud-scoring consumer. On a traditional queue-based broker, that history is very likely already gone, and standing up the same capability generally means adding a durable log or event-store layer on top of the broker.
Trade-offs and pitfalls
A traditional broker often has richer built-in per-message routing and prioritization features (complex routing keys, priority queues, per-message TTL (time-to-live)) that a commit-log platform doesn't naturally provide out of the box. If your use case is genuinely simple task distribution among workers, with no need for replay or multiple independent readers, a traditional broker's simpler operational model can still be the better fit; reaching for a commit-log platform by default adds retention and partitioning complexity you may not need.
Explain the core building blocks of Apache Kafka: topics, partitions, brokers, replication, and leader/follower roles for a partition. How does a producer's message end up durably stored and available to consumers?
Sample Answer
Direct answer
A producer's message becomes durable when it's written to the leader replica of a partition and then copied to enough in-sync replicas to satisfy the acknowledgment level the producer requested; consumers then read it in the exact order it was appended within that partition.
Structured elaboration
A topic is a named, logical stream of records, split into one or more partitions for parallelism. Each partition is an append-only, ordered log; every record gets a monotonically increasing offset within its partition. Each partition has one leader broker (which serves all reads and writes for it) and zero or more follower replicas that continuously copy the leader's log. Only followers that are caught up within a bounded lag are considered in-sync replicas (ISR); a write is durable once it's replicated across the ISR set, not merely written to the leader's local disk.
Worked example
Concretely: a producer sends a record with key user-42 to topic clicks. The partitioner hashes the key to pick, say, partition 3. The leader broker for partition 3 appends the record at the next offset (say offset 8842), then two followers replicate it. Once both are confirmed in-sync, the broker acknowledges the write back to the producer (assuming acks=all). A consumer subscribed to partition 3 will always see offset 8842 after offset 8841 and before offset 8843, because ordering is only ever guaranteed within a single partition, not across the whole topic.
Trade-offs and pitfalls
Because ordering is per-partition, spreading a topic across more partitions increases throughput and parallel consumption but only preserves ordering for records sharing the same partition (typically enforced via a consistent partition key). A common misconception is expecting global, topic-wide ordering by default; that only holds if the topic has exactly one partition, which caps throughput to a single broker's capacity for that topic.
Unlock Full Question Bank
Get access to all 9 Stream Processing and Event Streaming interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.