11 min readArchitecture
@Transactional Won't Save Your Kafka Publish
The most dangerous bug in a payment system is not the one that crashes. It is the one that quietly loses a transaction — no stack trace, no alert, nothing on a dashboard. Just a customer whose account was debited and whose money went nowhere.
This is a walk through an architecture built around two hard guarantees — zero message loss and zero duplicate processing — at high throughput. The mechanism at its centre is the transactional outbox, hash-partitioned across Oracle and drained into Kafka.
The real problem: dual writes
It starts with code that looks entirely reasonable. A payment request arrives, you persist it, you publish an event:
@Transactional
public void acceptPayment(PaymentRequest req) {
paymentRepository.save(new PaymentTransaction(req)); // Oracle
kafkaTemplate.send("payment.events.v1", new PaymentAccepted(req.id())); // Kafka
}
This is wrong, and it is wrong despite the @Transactional annotation.
Kafka is not enrolled in the database transaction. There are two writes to two systems with no atomicity between them, and three ways it fails:
- The database commits and the Kafka publish fails — a transaction with no event. Money taken, never processed.
- Kafka accepts the event and the database commit fails — an event with no transaction. Consumers process something that does not exist.
- The JVM dies between the two, which produces the first case.
Teams often treat this as a race condition and reach for retries. It is not a race condition. It is a correctness gap: you can narrow the window, you cannot close it.
The transactional outbox
The fix is almost boring. Instead of publishing the event, write it to a table in the same database, in the same transaction as the business data. Now either both commit or neither does, and Oracle guarantees that for you.
@Transactional
public void acceptPayment(PaymentRequest req) {
PaymentTransaction txn = paymentRepository.save(new PaymentTransaction(req));
outboxRepository.save(OutboxEvent.of(txn)); // same transaction, same database
}
A separate process — the relay — then reads the outbox table and publishes to Kafka. That separation is the point: the accept path only ever writes to the database, which is fast and reliable, and publishing becomes asynchronous.
The Oracle table, and why hash partitioning
At high throughput a single outbox table becomes a hotspot fast. Every insert lands at the same end of the table, contending on the same index blocks — what Oracle people call right-hand index growth.
Hash partitioning spreads that out:
CREATE TABLE outbox_event (
id NUMBER GENERATED ALWAYS AS IDENTITY,
transaction_id RAW(16) NOT NULL,
partition_id NUMBER(3) NOT NULL, -- what a relay claims
event_type VARCHAR2(64) NOT NULL,
payload CLOB NOT NULL,
created_at TIMESTAMP(6) DEFAULT SYSTIMESTAMP NOT NULL,
published_at TIMESTAMP(6),
CONSTRAINT pk_outbox PRIMARY KEY (id, transaction_id)
)
PARTITION BY HASH (transaction_id) PARTITIONS 32;
-- LOCAL: each partition gets its own index, so there is no shared structure to contend on
CREATE INDEX ix_outbox_unpublished
ON outbox_event (partition_id, published_at, id) LOCAL;
Hashing on transaction_id distributes inserts evenly across 32 partitions. No single partition takes the load, and because each has its own local index, index maintenance parallelises too.
Note the index is LOCAL, and that matters more than it looks. A GLOBAL index would put every partition back onto one shared structure and undo the entire benefit — you would have paid for partitioning and kept the hotspot.
The relay: claiming rows with SKIP LOCKED
Here is the interesting part. Multiple relay instances run concurrently, and no two may take the same row — but they must not wait on each other either.
Oracle’s FOR UPDATE SKIP LOCKED is built for exactly this. It skips rows another session has locked and takes the next available ones:
SELECT id, transaction_id, event_type, payload
FROM outbox_event
WHERE partition_id IN (?, ?, ?, ?) -- the partitions this relay holds
AND published_at IS NULL
ORDER BY id
FETCH FIRST 500 ROWS ONLY
FOR UPDATE SKIP LOCKED
Without SKIP LOCKED, relays block on each other’s locks and the whole thing degrades to effectively single-threaded. That one keyword is what makes the parallelism real.
The relay loop in Java:
@Component
@RequiredArgsConstructor
public class OutboxRelay {
private final JdbcTemplate jdbc;
private final KafkaTemplate<String, byte[]> kafka;
private final PartitionLease lease; // which partitions are mine — see below
private static final int BATCH = 500;
@Scheduled(fixedDelay = 20)
@Transactional
public void drain() {
List<Integer> mine = lease.heldPartitions();
if (mine.isEmpty()) return; // holding nothing yet
List<OutboxRow> rows = jdbc.query(
claimSql(mine.size()), rowMapper, mine.toArray());
if (rows.isEmpty()) return;
List<CompletableFuture<?>> sends = new ArrayList<>(rows.size());
for (OutboxRow row : rows) {
ProducerRecord<String, byte[]> record = new ProducerRecord<>(
"payment.events.v1",
null,
row.transactionId(), // key = transaction id, so ordering holds
row.payload());
// the consumer dedupes on this
record.headers().add("event-id", Long.toString(row.id()).getBytes(UTF_8));
sends.add(kafka.send(record));
}
// wait for every ack BEFORE marking anything published
CompletableFuture.allOf(sends.toArray(CompletableFuture[]::new))
.join();
jdbc.batchUpdate(
"UPDATE outbox_event SET published_at = SYSTIMESTAMP WHERE id = ?",
rows.stream().map(r -> new Object[]{ r.id() }).toList());
}
}
Two decisions in there are deliberate.
published_at is only set after Kafka acknowledges. If the relay dies after publishing but before updating, the transaction rolls back, the locks release, and another relay picks the same rows up again. The event goes out twice. That is at-least-once delivery, and it is what we want — the alternative is at-most-once, which loses messages. In payments a duplicate is far cheaper than a loss, because a duplicate can be absorbed at the consumer and a loss cannot be recovered at all.
The Kafka key is transaction_id. Every event for one transaction lands on the same Kafka partition, so ordering is preserved where it matters. There is no global ordering across the topic — and no need for it. Per-transaction ordering is the actual requirement.
Deciding which relay owns which partitions
Relays need to agree on who drains what. A lease table handles it:
CREATE TABLE relay_lease (
partition_id NUMBER(3) PRIMARY KEY,
owner_id VARCHAR2(64),
expires_at TIMESTAMP(6)
);
Each relay renews its own leases on a heartbeat. When one dies its leases expire and the others claim them:
@Scheduled(fixedDelay = 5_000)
@Transactional
public void renewAndClaim() {
jdbc.update("""
UPDATE relay_lease
SET expires_at = SYSTIMESTAMP + INTERVAL '30' SECOND
WHERE owner_id = ?
""", ownerId);
// claim expired leases — SKIP LOCKED again, so relays do not fight over them
jdbc.update("""
UPDATE relay_lease
SET owner_id = ?, expires_at = SYSTIMESTAMP + INTERVAL '30' SECOND
WHERE partition_id IN (
SELECT partition_id FROM relay_lease
WHERE expires_at < SYSTIMESTAMP
FETCH FIRST ? ROWS ONLY
FOR UPDATE SKIP LOCKED)
""", ownerId, targetPartitionCount());
}
Keep the lease duration comfortably longer than the heartbeat interval — 5 seconds against 30 here. Set them too close and an ordinary GC pause or a moment of network jitter hands partitions back and forth for no reason, disrupting ordering while nothing is actually wrong.
Kafka configuration
A few producer settings are not up for discussion. Without them the guarantees above collapse:
spring:
kafka:
producer:
acks: all # no ack until every in-sync replica has it
properties:
enable.idempotence: true # suppresses duplicates from producer retries
max.in.flight.requests.per.connection: 5
compression.type: lz4
batch-size: 65536
properties.linger.ms: 10 # a little latency, a lot of throughput
acks=all with replication factor 3 means losing a broker does not lose messages. enable.idempotence=true removes duplicates caused by the producer’s own retries — but note it does nothing about the duplicates the relay can create in the crash scenario above. That is the consumer’s job.
linger.ms: 10 lets the producer wait briefly to batch more records. At high throughput the difference is large, but it is a conscious trade against your latency budget, not a free win.
The consumer: idempotency is the real safeguard
Since delivery is at-least-once, consumers must tolerate duplicates — and “tolerate” means not charging someone twice, not merely logging that it happened.
The reliable way is to let a database constraint decide, not an if:
@KafkaListener(topics = "payment.events.v1", groupId = "payment-processor")
@Transactional
public void onPaymentAccepted(ConsumerRecord<String, byte[]> record) {
String eventId = header(record, "event-id");
try {
// the UNIQUE constraint is the guarantee; a code check is only an optimisation
processedEvents.insert(eventId);
} catch (DuplicateKeyException e) {
log.debug("Ignoring duplicate event {}", eventId);
return;
}
PaymentAccepted event = deserialize(record.value());
paymentProcessor.process(event); // advance the state machine
}
processed_events has a UNIQUE constraint on event_id, and the insert happens in the same transaction as the processing. If processing fails, the idempotency record rolls back with it and the event is retried. Put them in separate transactions and you eventually get an event marked “processed” whose work never happened — the worst outcome of the three.
A plain SELECT check (“have I seen this?”) is not sufficient. Two consumers can check concurrently, both see nothing, and both proceed. Only the constraint actually answers the question.
Java 21 and virtual threads
Virtual threads genuinely fit here. The accept path is I/O-bound — waiting on the database, waiting on a bank. On platform threads each concurrent request pins an OS thread, and a few thousand of those is the ceiling.
In Spring Boot 3.2+ it is one line:
spring:
threads:
virtual:
enabled: true
Two caveats that get skipped.
Virtual threads do not enlarge your connection pool. A hundred thousand virtual threads will queue for the same ten database connections. They remove thread scarcity, not connection scarcity. Measure where the bottleneck actually is, or you have only moved the waiting somewhere less visible.
Blocking inside synchronized pins the carrier thread on JDK 21, which defeats the point. Older libraries still do this; -Djdk.tracePinnedThreads=full will find it.
Protecting what is downstream
The most realistic thing on the whole diagram is a small note in the corner: bank calls are rate limited, perhaps to 10K TPS.
That is the crux. You can accept a million transactions per second; the bank or MFS provider cannot. They will take a few thousand. The system’s real job is absorbing that gap in a controlled way — Kafka acts as the buffer, and consumers run at whatever pace downstream allows.
Resilience4j expresses this declaratively:
@RateLimiter(name = "bankConnector") // token bucket matched to the bank's limit
@Bulkhead(name = "bankConnector") // cap concurrent calls
@CircuitBreaker(name = "bankConnector", fallbackMethod = "queueForRetry")
@Retry(name = "bankConnector") // exponential backoff
public BankResponse submit(PaymentTransaction txn) {
return bankClient.submit(txn);
}
private BankResponse queueForRetry(PaymentTransaction txn, Throwable cause) {
// the bank is down — park it for later rather than failing the transaction
retryQueue.enqueue(txn, cause);
return BankResponse.pending(txn.id());
}
The circuit breaker’s role here is subtle and important. When a downstream system slows down, retrying makes it worse — you are adding load to something already struggling. Failing fast gives it room to recover.
And note the fallback: when the bank is unavailable the transaction is parked, not lost. Let that backpressure reach Kafka. If consumers cannot make progress, consumer lag rises — and that is the single most important alert in this design.
Being honest about the numbers
Now the part that architecture diagrams tend to leave off.
The design targets 1M TPS against a single Oracle instance, with seven nines of availability. Both claims deserve scrutiny, because senior reviewers will apply it whether or not you do.
A single instance and seven nines are mutually exclusive. Seven nines is roughly three seconds of downtime per year. One database instance — however large — needs patching, needs restarting, and its hardware eventually fails. Availability is bounded by the weakest single point of failure, and here that is the database. Realistically this needs Oracle RAC, Data Guard, or at minimum a tested failover plan — at which point the “simpler than sharding” argument for a single instance has largely evaporated.
A million inserts per second on one instance is extremely aggressive. Partitioning genuinely helps, but the limit is not partition count — it is redo log throughput, the commit rate that LGWR can sustain, and I/O bandwidth. Every durable commit must write redo. The biggest lever is usually batching: instead of one commit per request, insert many rows per round trip.
jdbc.batchUpdate("INSERT INTO outbox_event (...) VALUES (?, ?, ?, ?)", batch);
That single change can cut commits per second by an order of magnitude, and commits are frequently the real ceiling.
My advice: treat these figures as targets, not commitments. The patterns — outbox, SKIP LOCKED relays, partition leases, idempotent consumers — are sound and well proven. The achievable throughput depends entirely on your hardware, redo configuration and transaction size. Load test it; do not assume it.
What actually matters
- Dual writes are never safe. Use an outbox so business data and event commit together.
- Hash partitioning removes the insert hotspot — and the index must be
LOCAL. FOR UPDATE SKIP LOCKEDis what lets relays work in parallel instead of queueing behind each other.- Mark
published_atonly after the Kafka ack, keeping delivery at-least-once rather than at-most-once. - Enforce consumer idempotency with a constraint, not an
if, and in the same transaction as the work. - Downstream capacity is the real limit. Buffer in Kafka; protect the bank with rate limiting, bulkheads and a circuit breaker.
- Challenge the numbers, particularly any single-instance database claim.
One closing thought. None of this complexity exists to make the system fast. It exists so that no transaction is ever lost. There are far simpler ways to reach high throughput if you are willing to drop a message occasionally — payments is precisely the domain where you are not, and every awkward piece of this design traces back to that.