04 / 05

Enumerate the failure windows between consuming an event, updating a database and committing the offset.

Difficulty: 8/10
Distributed-system reasoning, Ordering, Exactly-once

Failure Windows Between Consuming, Updating the Database, and Committing the Offset

There are four failure windows between consuming an event, updating a database, and committing the offset. The first window is between the poll and the database update: the consumer has read the event but has not yet written to the database. If the consumer crashes here, the event is not processed, the offset is not committed, and the event will be re-delivered on restart. This is safe: no data loss, no duplicate side effect. The second window is between the database update and the offset commit: the consumer has written to the database but has not committed the offset. If the consumer crashes here, the database has the update but the offset is not advanced. On restart, the event is re-delivered and the database update is applied again. This is the classic duplicate window: the side effect happens twice unless the consumer is idempotent. The third window is during the offset commit itself: the consumer sends the commit but crashes before receiving the acknowledgment. The commit may or may not have been persisted. On restart, the consumer may resume from the old offset (if the commit was not persisted) or from the new offset (if it was). If it resumes from the old offset, the event is re-delivered and the database update is applied again, causing a duplicate. The fourth window is after the offset commit: the consumer has committed the offset and the database update is durable. If the consumer crashes here, no re-delivery occurs. This is the safe state. The trade-off is between the order of operations and the type of failure. Committing before the database update gives at-most-once (potential data loss); committing after gives at-least-once (potential duplicates). The standard choice is to commit after the database update and to make the consumer idempotent.

The mechanism for each window is determined by the order of operations and the durability of each step. If the database update is in a transaction, the transaction must commit before the offset commit; if the transaction rolls back, the offset should not be committed. If the consumer crashes between the database commit and the offset commit, the database has the change but the offset is old, so the event is re-delivered. The consumer must detect that the event was already processed. This is where the inbox pattern or idempotent operations are needed. The trade-off is between the cost of idempotency and the risk of duplicates. For a system where duplicates are unacceptable, the consumer must use a deduplication store or idempotent operations. For a system where duplicates are tolerable, at-least-once is sufficient. Version note: Kafka transactions can make the database update and the offset commit atomic if the database update is also a Kafka write (consume-transform-produce). But if the database is external, there is no shared transaction, and the duplicate window remains. The practical approach is to accept at-least-once and make the consumer idempotent. This is the standard pattern in event-driven systems. Also note that auto-commit changes the windows: with auto-commit, the commit can happen before the database update, which gives at-most-once and potential data loss. This is why manual commit is recommended for consumers that write to external systems.

A common mistake is to commit the offset before the database update, which gives at-most-once and can lose data. Another mistake is to assume that the database update and the offset commit are atomic; they are not, unless the database is Kafka. A third mistake is to use auto-commit without understanding its behavior. The trade-off is between the cost of idempotency and the risk of duplicates or loss. A good practice is to commit the offset after the database transaction commits, and to make the database update idempotent using the event ID. This gives at-least-once delivery with effectively-once side effects. If the database update fails, the offset is not committed, and the event is re-delivered. If the database update succeeds but the offset commit fails, the event is re-delivered and the idempotent update is a no-op. Version note: the inbox pattern is the standard way to implement this. The inbox table has a unique constraint on the event ID; the insert and the business update are in the same transaction. If the insert fails with a duplicate key, the business update is skipped. This gives idempotency without a separate deduplication check.

javascript
  1. 1

    Four windows: before database update, between update and commit, during commit, after commit.

  2. 2

    Window 1: crash before update -> re-delivered, safe.

  3. 3

    Window 2: crash after update before commit -> re-delivered, duplicate unless idempotent.

  4. 4

    Window 3: crash during commit -> may re-deliver or not, duplicate possible.

  5. 5

    Window 4: crash after commit -> no re-delivery, safe.

  6. 6

    Commit after the database update for at-least-once; commit before for at-most-once.

  7. 7

    Make the consumer idempotent with the inbox pattern to handle window 2 and 3.

  8. 8

    Kafka transactions only help if the output is also Kafka.

Share

Share via WhatsApp, X, Facebook, LinkedIn or copy link. Open Graph preview enabled.