Making a Consumer Idempotent for Safe Replay
A consumer is idempotent if processing the same event more than once produces the same result as processing it once. This is essential for safe replay because Kafka delivers at least once, and rebalances, retries, and DLT replays can all cause duplicate delivery. There are three main techniques: event IDs with a deduplication store, idempotent operations, and uniqueness constraints in the target system. Event IDs with a deduplication store: every event carries a unique ID, and the consumer checks whether it has already processed that ID before applying the effect. The deduplication store can be a database table with a unique constraint, a Redis key with a TTL, or a compacted Kafka topic. Idempotent operations: instead of 'increment counter by 1', use 'set counter to value' or 'upsert row with version'. This makes the operation naturally idempotent without a separate deduplication store. Uniqueness constraints: the target system has a primary key or unique index that rejects duplicate inserts. This is a form of deduplication at the data layer. The trade-off is between storage cost and simplicity. A deduplication store grows over time and needs cleanup; idempotent operations are simpler but not always possible; uniqueness constraints require the target system to support them. Version note: Kafka's exactly-once semantics can reduce duplicates for consume-transform-produce within Kafka, but they do not eliminate the need for idempotency when writing to external systems. The idempotency is always the consumer's responsibility.
The mechanism for the deduplication store is to write the event ID and the effect in the same transaction. For a database-backed consumer, the pattern is: begin transaction, insert the event ID into an inbox table (which fails if the ID already exists), apply the effect, commit. If the insert fails with a duplicate key, the consumer skips the effect and commits or rolls back without harm. This is the inbox pattern. The mechanism for idempotent operations is to design the update so that applying it twice has the same result. For example, instead of 'balance = balance + amount', use 'balance = new_balance' where new_balance is computed from the event. This requires the event to carry the full new state, not just the delta. The mechanism for uniqueness constraints is to rely on the database's primary key or unique index; if the consumer tries to insert a duplicate, the database rejects it. The trade-off is between the richness of the event and the complexity of the consumer. An event that carries the full state is easier to apply idempotently but larger; an event that carries only the delta requires a deduplication store. Version note: the inbox pattern is a standard integration pattern and is often implemented with a database table. Some frameworks provide idempotent consumer support, such as Spring Kafka's idempotent receiver or Kafka Streams' exactly-once processing. Check the framework version and features.
A common mistake is to use an in-memory set of seen event IDs. This fails on restart and does not work across consumer instances. Another mistake is to check for duplicates and then write without a transaction, which leaves a race window where two instances can both check and both write. A third mistake is to assume that the consumer will never see duplicates, which is false in an at-least-once system. The trade-off is between the cost of deduplication and the risk of corruption. For a financial system, the cost of deduplication is worth it; for a metrics system, it may not be. The key is to understand the business impact of duplicates and to choose the right level of protection. Version note: if you use Kafka transactions for consume-transform-produce, the consumer offsets and the output records are committed atomically, which eliminates duplicates within Kafka. But if the consumer writes to an external system, you still need idempotency. For replay, you need idempotency because replay re-delivers events that were already processed.
Idempotent means processing the same event twice has the same effect as processing it once.
Use event IDs plus a deduplication store (inbox table, Redis, compacted topic).
Prefer idempotent operations (set, upsert) over non-idempotent ones (increment, append).
Write the event ID and the effect in the same transaction to avoid race windows.
In-memory deduplication fails on restart and across instances.
Uniqueness constraints in the target system provide deduplication.
Kafka exactly-once helps within Kafka but not for external systems.
Replay requires idempotency because events are re-delivered.
0-2 years experience
2-5 years experience
5-8 years experience
8+ years experience