Design a Kafka Consumer Service That Groups Trace Events by Trace ID
Company: New Relic
Role: Backend Engineer
Category: System Design
Difficulty: medium
Interview Round: Onsite
Design a distributed event-processing system that consumes trace events from Kafka and groups together all events that belong to the same trace, so that each trace can be processed as a whole.
### Constraints and Clarifications
- Events are received through Kafka.
- Every event contains a `traceId`, and all events of one trace share the same `traceId`.
- A single trace can last anywhere from a few seconds to several hours.
- The processing logic itself is out of scope. Assume a function already exists that processes all events of one trace together; this question calls it `process_trace(trace_id, events)`.
- No event rates, trace sizes or latency targets are specified, so part of the exercise is to ask for them or to state them as assumptions.
### Clarifying Questions
- What event throughput and how many concurrently open traces must the system handle, and how large can a single trace become?
- Is the Kafka topic already keyed by `traceId`, or by something else such as the emitting service?
- Does any event mark the end of a trace (for example the completion of its root span), or must completion be inferred?
- How far out of order, and how late, can the events of one trace arrive?
- Must each trace be processed exactly once, and is `process_trace` safe to call twice for the same trace?
- How soon after a trace ends must it be processed?
### Part 1 — Group events by traceId
Explain how all events of one trace end up in the same place. Start with the simplest approach, buffering events in memory inside the Kafka consumers, and then assess whether that approach scales.
```hint Who owns a trace
Kafka assigns each partition to exactly one consumer in a group at a time. Think about what decides which partition an event lands in, and what happens to buffered state when that assignment changes.
```
#### What This Part Should Cover
- How the events of one trace are routed to a single consumer, and what to do if the topic is not keyed that way.
- The in-memory design and its limits: memory growth, consumer restarts and group rebalances, and when Kafka offsets can safely be committed.
- A clear verdict on whether in-memory grouping scales when traces last hours.
### Part 2 — Keep trace state in an external datastore
Move per-trace state out of consumer memory into an external datastore. Describe the data model, the write path for each incoming event, how the design stays within memory limits when traces run for hours, and how the grouped events of a trace are retrieved when it is time to process it.
```hint Two kinds of state
The information needed to decide whether a trace is finished is small, while the events themselves can be large. Consider whether they belong in the same place and under the same access pattern.
```
#### What This Part Should Cover
- The choice of store and a schema keyed by `traceId`, with writes that tolerate duplicate delivery.
- How consumer memory stays bounded regardless of trace duration, and how Kafka offset commits relate to durable writes.
- The retrieval path for processing, including traces too large to load at once.
- The costs this adds compared with the in-memory design.
### Part 3 — Decide that a trace is complete and process it
Decide when a trace is finished, trigger `process_trace` for it, and make the design work for traces that run for hours as well as for traces that last seconds.
```hint One rule rarely fits both
A rule tuned for traces that last seconds can misfire on traces that last hours. Consider which signals you can combine, and what happens when an event arrives after you declared its trace complete.
```
#### Clarifying Questions for this Part
- Must a trace that runs for hours be processed only once, at the end, or are periodic partial results acceptable?
- What should happen to an event that arrives after its trace has already been processed?
#### What This Part Should Cover
- Completion signals and their failure modes: an explicit end event, an inactivity timeout, a maximum duration.
- A way to find traces that are due for completion without scanning every open trace.
- Claiming a trace so that it is processed once, and handling retries and late events.
- A specific treatment for traces that stay open for hours.
### What a Strong Answer Covers
- Stated assumptions for throughput, trace size and latency, since the prompt gives none, and design choices that follow from them.
- A comparison of in-memory buffering, an external datastore and a stream processor with managed keyed state, with explicit trade-offs.
- Failure handling for consumer crashes, rebalances, duplicate delivery and datastore outages.
- Observability for open traces, completion delays, late events and consumer lag.
### Follow-up Questions
- A few traces are orders of magnitude larger than the rest. How does your design avoid hot partitions in Kafka and hot keys in the datastore?
- The topic is replayed from an earlier offset after a bug fix. How do you avoid processing traces twice, or processing them with only part of their events?
- How would the design change if some processing had to happen while a trace is still open, for example alerting on an error within seconds?
- If the datastore is unavailable for several minutes, what happens to consumption and to traces that were about to complete?
Overview: Design a distributed service that consumes trace events from Kafka, groups every event sharing a trace ID, and hands each complete trace to an existing processing function. Covers partitioning, in-memory versus external state, deciding when a trace is complete, and traces that stay open for hours.
Read the full New Relic Backend Engineer interview experience this question came from