Kafka Design Patterns Cheat Sheet: Guarantees and Tradeoffs
Recall Kafka ordering, transactions, replay, stateful processing and integration patterns through concise comparisons and concrete design checks.
Use this reference after studying a concept, then apply it to a fresh practice question . The useful interview answer names a boundary and its consequences.
Identify what the message means
| Construction | Meaning | Check before choosing it |
|---|---|---|
| Command | A request for an owner to perform work; rejection is possible. | Who owns the action, and what response establishes its outcome? |
| Event | A fact about something that happened. | Does its identity survive retries, and can consumers interpret it independently? |
| Document/state update | A representation of current or versioned state. | Is it a complete replacement or a patch? What does a tombstone mean? |
| Claim check | A message carries a reference to separately stored content. | Will authorized replay readers still have the referenced object and its integrity metadata? |
Name the guarantee boundary
| Mechanism | What it helps establish | What still needs design |
|---|---|---|
| Stable key placement | Related records can share one partition’s append order. | Business sequence, concurrent writers and partition-count changes. |
| Producer idempotence | Certain producer retry duplicates are prevented. | Repeated business events and effects outside Kafka. |
| Kafka transactions | Atomic visibility of a set of Kafka writes and participating offsets to appropriate readers. | Database writes, HTTP side effects and consumers using unsuitable isolation. |
| Idempotent consumer | Repeated delivery does not repeat a defined business effect. | Stable identity, atomic deduplication with the effect, and ledger retention. |
| Transactional outbox | A database transaction records business state and the intent to publish together. | Relay retries, downstream deduplication, ordering and cleanup. |
| Compaction | Retained records converge toward latest keyed state subject to compaction behavior. | Complete event history, immutable audit and time-bounded deletion. |
For example, committing an input offset after sending an HTTP request does not make that HTTP effect atomic with Kafka. A timeout leaves an uncertain result. Use an idempotency contract or authoritative reconciliation before retrying an irreversible effect.
Trace state instead of recognizing a pattern name
- A KStream represents record occurrences. Filtering an update out does not automatically retract an earlier output occurrence.
- A KTable represents keyed updates. A row that no longer satisfies a table filter can retract the prior result.
- A stream-table join looks up table state as stream records are processed; later table updates do not ordinarily re-emit past stream records.
- A stream-stream join matches occurrences within its time bounds. Account for both input rates, retention and multiple matches.
- Event time assigns meaning to timestamps; stream time advances from processed timestamps. An idle input does not supply a wall-clock expiry mechanism.
- A grace period governs accepted lateness for a window operation. Check the particular window type and API; session merging adds its own inactivity-gap semantics.
Read the Kafka Streams DSL guide for exact API behavior and test a minimal trace on your deployed version.
Make replay a bounded operation
Before approving replay, identify the baseline state, inclusive or exclusive source boundary, retained records, schemas, code/configuration version and external-effect policy. Resetting consumer offsets alone does not restore a consistent application snapshot.
A sink’s durable frontier should mean that every record through that partition offset has been applied. A frontier of 95 does not prove that no higher offset was individually applied. Compare sinks at a common offset vector, including key sets, revisions and values; matching totals can conceal compensating errors.
Calculate with explicit assumptions
| Estimate | Relationship | Qualification |
|---|---|---|
| Backlog drain time | Backlog / (processing rate − continuing arrival rate) | Processing must exceed arrivals; rates need compatible units and sustained capacity. |
| Retained record count | Combined input rate × retention duration | Include both join inputs when both are stored; this is before implementation overhead. |
| Payload-state size | Concurrent entries × bytes per entry | Add keys, indexes, serialization and storage overhead separately. |
| Parallel prerequisite latency | Approximately the slowest prerequisite, plus orchestration overhead | Dependencies must actually be parallel; retries and tail latency need separate treatment. |
Separate a guarantee from a test result
| When reviewing… | Ask… |
|---|---|
| A replay rule | Is the required result the historical decision or a reinterpretation under current policy? Where is the chosen rule version recorded? |
| A stateful join | Which input triggers output? What happens to a result already emitted when the lookup state changes? |
| A recovery rehearsal | Did assignment, decoding, state correctness or throughput fail? Which controlled comparison isolates the remaining bottleneck? |
| A capacity estimate | Are the rates measured under the same workload, and does continuing traffic leave positive recovery capacity? |
After choosing an answer, explain why each alternative fails under the stated conditions. Then change one condition that would make a different design appropriate.
Use a failure to test the design
Ask what happens if a worker fails after the effect but before the acknowledgement, a producer retries with a new identity, the partition count changes, a schema reader lags behind, or a replay passes a newer live update. Then describe the observable evidence that distinguishes safe recovery from silent loss or duplication.
Practise these decisions in IT Mastery and send a correction if a case needs clearer assumptions.