A team bridges domain events to Kafka with @TransactionalEventListener(AFTER_COMMIT). Occasionally the DB shows a change but no Kafka message was produced. Diagnose and design a fix.
answer
- two resources, no shared tx = dual-write
- crash/broker-down after commit -> lost event
- outbox row inserted in SAME tx = atomic
- relay: polling or Debezium CDC, retries
- at-least-once -> idempotent consumers by event id
basics
~20 sAFTER_COMMIT runs only after the DB commits, but the Kafka send is a separate operation. If the app crashes or the broker is unavailable after commit but before the send succeeds, the message is lost. This dual-write gap is fixed with a transactional outbox.
solid answer
~50 sThe symptom is the classic dual-write problem. `@TransactionalEventListener(AFTER_COMMIT)` correctly waits for the DB commit, but the commit and the `KafkaTemplate.send` are two independent resources with no shared transaction. Between them the process can crash, the broker can be down, or the send can fail asynchronously — the row is durable, the event never leaves. AFTER_COMMIT explicitly cannot make them atomic. The standard fix is the **transactional outbox**: within the same DB transaction that changes the aggregate, insert a row into an `outbox` table. A separate relay — a polling publisher or CDC (Debezium) reading the DB log — reads unsent outbox rows and publishes them to Kafka, marking them sent, retrying on failure. This gives at-least-once delivery keyed to the committed DB state; consumers must be idempotent (dedupe by event id). It trades a bit of latency and a background process for eliminating lost events.
code
java · 36 lines// Outbox write happens in the SAME transaction as the business change.
@Service
class OrderService {
private final OrderRepo orders;
private final OutboxRepo outbox;
private final ObjectMapper json;
OrderService(OrderRepo o, OutboxRepo ob, ObjectMapper j){ orders=o; outbox=ob; json=j; }
@Transactional
public void place(Order order) throws JsonProcessingException {
orders.save(order); // (1) aggregate
outbox.save(new OutboxRecord( // (2) event, same tx -> atomic
UUID.randomUUID(),
"OrderPlaced",
order.getId(),
json.writeValueAsString(new OrderPlaced(order.getId())),
OutboxStatus.NEW));
}
}
// Separate relay ships committed outbox rows to Kafka and retries.
@Component
class OutboxRelay {
private final OutboxRepo outbox;
private final KafkaTemplate<String, String> kafka;
OutboxRelay(OutboxRepo o, KafkaTemplate<String,String> k){ outbox=o; kafka=k; }
@Scheduled(fixedDelay = 500)
@Transactional
public void flush() {
for (OutboxRecord r : outbox.findTop100ByStatusOrderByCreatedAt(OutboxStatus.NEW)) {
kafka.send("orders", r.getAggregateId().toString(), r.getPayload()); // key = order id -> ordering
r.markSent(); // at-least-once: a crash before this replays the row
}
}
}go deeper
Out of depth; only needs to sense that DB and broker can disagree.
Can name the dual-write problem and that AFTER_COMMIT isn't atomic.
Can prescribe the outbox and explain at-least-once + idempotency.
Owns the full design: outbox vs CDC vs XA trade-offs, ordering, dedupe, latency, and how the in-JVM event still fits.
## Diagnosis: the dual-write problem The pipeline does **two writes to two systems**: (1) commit the aggregate to the database, (2) send a message to Kafka. `@TransactionalEventListener(AFTER_COMMIT)` sequences them so the send happens **after** commit — good, because you never announce uncommitted data. But it does **not** make them **atomic**. Failure windows: - **Crash after commit, before send** — JVM dies once the transaction has committed but before/while calling `send`. Row persisted, no message. - **Broker unavailable** — Kafka is down at send time; the send fails and there's no transaction to roll the DB back (it's already committed). - **Async send failure** — `KafkaTemplate.send` returns a future; if you don't block/handle the callback, a broker-side failure is silently dropped. - **REQUIRES_NEW confusion** — writes done in the listener may or may not persist depending on propagation, compounding inconsistency. This is inherent: two resources, no common commit coordinator (short of rarely-used XA/2PC, which brokers like Kafka don't support well). ## Fix: transactional outbox **Single local transaction, single resource.** In the same DB transaction that mutates the aggregate, also `INSERT` the event into an **`outbox`** table (id, aggregate id, type, payload, created_at, status). Because it's the same transaction, the outbox row is committed **atomically** with the business change — no gap. **Relay to the broker** (separate process/thread): - **Polling publisher** — periodically `SELECT` unsent rows, publish to Kafka, mark `sent` (or delete). Simple, works anywhere. - **CDC (Debezium)** — tail the database transaction log and stream inserted outbox rows to Kafka. No polling, lower latency, but more infrastructure. The relay **retries** until the broker acks, so delivery becomes **at-least-once** anchored to committed state. A crash just means the relay re-reads the row later. ## Consequences you must design for - **At-least-once ⇒ idempotent consumers.** The same outbox row may be published more than once (relay crash after send, before marking sent). Consumers dedupe by a stable **event id**. - **Ordering.** Preserve per-aggregate order via the partition key (e.g. aggregate id) and monotonic outbox ordering. - **Outbox growth.** Purge/archive sent rows. - **Latency.** Polling adds a small delay; tune interval or use CDC. ## Where the in-JVM event still fits You can keep publishing an in-JVM `@EventListener`/`@TransactionalEventListener` for **local** side effects, but the **cross-process** hop should go through the outbox, not a bare AFTER_COMMIT `send`. So the bridge becomes: transaction writes aggregate + outbox row → relay → Kafka. ## Alternatives (and why usually not) - **XA / 2PC across DB and broker** — Kafka lacks solid XA support; heavy, poor performance, operational pain. - **Listener-based publish with retries only** — still loses events on crash between commit and any retry attempt; no persistence of intent. The outbox is the pragmatic, widely-adopted answer.
- Why not use XA/2PC to make the DB commit and Kafka send atomic instead?Kafka has no robust XA support and 2PC is operationally heavy with poor performance and blocking failure modes. The outbox achieves practical atomicity using only the database's local transaction, which is simpler and more reliable.
- The outbox gives at-least-once. How do you prevent duplicate side effects downstream?Make consumers idempotent: dedupe on a stable event id (or use an idempotency key / processed-events table), and design handlers so reprocessing the same event yields the same state.
- How do you preserve event ordering with the outbox?Partition by the aggregate id (Kafka key = aggregate id) so all events for one aggregate land on the same partition, and have the relay publish outbox rows in creation order.
saying these in an interview costs you the question
- Insisting AFTER_COMMIT already guarantees delivery
- Proposing XA/2PC with Kafka as the primary fix
- Forgetting consumers must be idempotent under at-least-once
- Ignoring ordering/partition-key when moving to the broker
- Blindly trusting KafkaTemplate.send without handling the async result