Quarkus SmallRye Reactive Messaging and Stripe Integration: How @Incoming Consumer Re-invocation on Message Nack Generates UUID_B on Redelivery, Mutiny Chain Deferred UUID Evaluation Produces a New Key on Re-subscription, and @Outgoing Kafka Producer Inside @Transactional Commits Messages Before the Application Transaction Rolls Back Triggering Duplicate Consumer Executions
The previous posts in this series traced Stripe idempotency failures to specific framework mechanisms: AOP proxy ordering in Spring Boot, reactive cold publisher re-subscription in Spring Data MongoDB, CDI interceptor priority in Quarkus Hibernate Reactive, and Testcontainers testing blind spots in Spring Boot integration tests. This post shifts to Quarkus SmallRye Reactive Messaging: the messaging layer that connects microservices via Kafka, AMQP, or in-memory channels. Three failure modes are mechanically distinct from each other and from all prior posts. Each is rooted in a different aspect of how SmallRye Reactive Messaging interacts with @Transactional, Mutiny’s deferred evaluation semantics, and Kafka’s producer isolation from JTA transactions.
SmallRye Reactive Messaging is not a retry library like @Retryable. It is a message-driven framework where “retry” means message redelivery from the broker: the consumer nacks the message (explicitly or by throwing), and the broker puts it back on the queue or topic. Every redelivery calls the @Incoming consumer method from the beginning. Every call that begins with UUID.randomUUID() generates a new UUID. The pattern is different from AOP retry but the Stripe double-charge outcome is identical: ch_A and ch_B, one per delivery attempt.
The third failure mode is architecturally distinct from both: it does not involve a consumer retry at all. It involves the gap between a Kafka producer’s commit and the JPA transaction’s commit. When those two are not enrolled in the same XA transaction (which is the default in Quarkus applications), a JPA rollback leaves a Kafka message already in the broker. A downstream consumer processes the message, calls Stripe, and creates ch_A. The application-level retry (which does not know the Kafka message was already produced) re-invokes the service, produces another message, and the consumer creates ch_B.
Background: SmallRye Reactive Messaging acknowledgment model and the @Incoming consumer lifecycle
SmallRye Reactive Messaging delivers messages to @Incoming-annotated consumer methods and manages acknowledgment based on the method’s return type and the @Acknowledgment annotation. When using the default @Acknowledgment(Strategy.POST_PROCESSING) strategy (which applies when the annotation is omitted), SmallRye acks the message if the consumer method completes normally (or if the returned Uni<Void> completes successfully) and nacks the message if the consumer throws an exception (or if the returned Uni<Void> fails with an error).
A nack means the message is negatively acknowledged. The exact consequence of a nack depends on the connector. For Kafka, the consumer offset is not advanced — the next poll will re-deliver the message. For AMQP, the message is returned to the queue (or moved to a dead-letter queue, depending on configuration). For SmallRye’s in-memory channel, nacks trigger configurable retry behavior. In all cases, the consumer method is eventually called again with the same message payload.
When SmallRye calls the consumer method again, it is a new method invocation. The Java stack frame is fresh. All local variables are re-initialized. Every expression in the method body, including UUID.randomUUID(), is re-evaluated. This is the fundamental mechanism behind all three failure modes in this post: redelivery = method re-entry = new UUID.
The Stripe idempotency contract requires that all retry attempts for the same logical billing event carry the same Idempotency-Key header. If UUID.randomUUID() is called at the top of the consumer method, attempt 1 gets UUID_A and attempt 2 (via redelivery) gets UUID_B. Stripe treats UUID_A and UUID_B as two distinct billing operations and creates ch_A and ch_B. The customer is charged twice. This is identical to the @Retryable re-execution failure in Spring Boot, except the re-execution is driven by message redelivery from the broker rather than by Spring Retry’s AOP advice calling proceed().
Mode 1: @Incoming consumer with @Transactional — Stripe committed before the database write fails, SmallRye nacks the message, redelivery generates UUID_B
The developer writes a Quarkus @Incoming consumer method annotated with both @Incoming("billing-events") and @Transactional. The intent is to make the consumer’s database write and the Stripe API call atomic: if either fails, the transaction rolls back and the message is redelivered for a clean retry. This reasoning is sound for database operations. It is incorrect for Stripe.
Stripe’s HTTP API is not a JTA-enrollable resource. The Stripe Java client makes an HTTPS POST request to api.stripe.com. There is no two-phase commit protocol between Stripe and your JTA transaction manager. When the consumer calls stripe.charges().create(params) inside the @Transactional boundary, the Stripe charge is committed at Stripe’s servers the moment Stripe returns a 2xx response. If the consumer then writes a BillingAttempt row to the database and that write fails (constraint violation, connection loss, optimistic lock contention), the @Transactional interceptor rolls back the database transaction — but the Stripe charge ch_A is already committed and cannot be rolled back.
SmallRye receives the exception thrown by the consumer. Because the exception propagated out of the consumer method body, SmallRye nacks the message. The Kafka consumer offset is not advanced; the next poll re-delivers the same message. SmallRye invokes the consumer method again. At method entry, UUID.randomUUID() generates UUID_B. The consumer calls Stripe with UUID_B. Stripe creates ch_B. The customer is charged twice.
// BillingEventConsumer.java — unsafe Mode 1 consumer
// @Transactional gives a false sense of atomicity between Stripe and the DB write.
// Stripe is not a JTA resource. ch_A is committed at Stripe before the DB write is attempted.
// If the DB write fails, @Transactional rolls back the DB write — but not ch_A.
// SmallRye nacks the message. Redelivery: UUID_B → ch_B.
@ApplicationScoped
public class BillingEventConsumer {
@Inject StripeClient stripeClient;
@Inject BillingAttemptRepository billingAttemptRepository;
// Developer's reasoning:
// "@Transactional makes the consumer atomic.
// If the DB write fails, the whole consumer rolls back and the message is redelivered.
// The retry starts clean because the DB write was rolled back."
//
// The problem:
// Stripe.charges().create() is called BEFORE repository.persist().
// ch_A is committed at Stripe when Stripe returns 200.
// repository.persist() throws a ConstraintViolationException (duplicate key, etc.).
// @Transactional rolls back the Panache persist — but not the Stripe charge.
// SmallRye nacks the message.
// Kafka re-delivers the message.
// onBillingEvent() is called again from line 1.
// UUID.randomUUID() generates UUID_B.
// stripeClient.createCharge(UUID_B) → ch_B.
// Customer charged twice.
@Incoming("billing-events")
@Transactional
public void onBillingEvent(BillingEvent event) {
// UUID generated at method entry — re-generated on every consumer invocation.
String idempotencyKey = UUID.randomUUID().toString();
ChargeParams params = ChargeParams.builder()
.customerId(event.getStripeCustomerId())
.amountCents(event.getAmountCents())
.currency("usd")
.idempotencyKey(idempotencyKey)
.build();
// Stripe called first — ch_A committed at Stripe.
Charge charge = stripeClient.createCharge(params);
// DB write second — may throw ConstraintViolationException if customer
// already has a BillingAttempt for this billing period.
BillingAttempt attempt = new BillingAttempt(
event.getCustomerId(),
charge.getId(),
event.getBillingPeriod(),
idempotencyKey,
BillingStatus.SUCCEEDED);
billingAttemptRepository.persist(attempt); // may throw → @Transactional rolls back
// If persist() throws:
// - @Transactional rolls back: BillingAttempt row NOT written to DB.
// - Stripe charge ch_A: already committed. Cannot be rolled back.
// - SmallRye nacks "billing-events" message.
// - Kafka re-delivers: onBillingEvent() called again.
// - UUID.randomUUID() → UUID_B.
// - stripeClient.createCharge(UUID_B) → ch_B.
}
}
The failure is not caused by calling Stripe inside a transaction. The failure is caused by calling UUID.randomUUID() at the consumer’s method entry, which means every consumer invocation — including redeliveries — generates a distinct idempotency key. The @Transactional annotation controls the database rollback behavior; it does not control the idempotency key’s stability across multiple consumer invocations.
This failure mode is distinct from the Spring @Retryable + @Transactional AOP ordering failure. In that case, the retry is driven by Spring Retry’s AOP interceptor re-calling proceed() on the intercepted method. Here, the retry is driven by Kafka message redelivery: the consumer method is called again by SmallRye’s polling loop, not by a retry interceptor. The mechanism is different; the UUID re-generation outcome is the same.
There is an additional subtlety specific to Quarkus’s @Transactional CDI interceptor integration with SmallRye. The @Transactional interceptor and the SmallRye acknowledgment interceptor are both CDI interceptors applied to the consumer method. Their relative priority (controlled by @Priority) determines whether the transaction commits before or after SmallRye acknowledges the message. If the transaction commits first and then the acknowledgment fails, the message may be redelivered even though the database write succeeded — leading to a duplicate processing scenario where the idempotency key guard in the database is the only protection. If the acknowledgment is sent first and the transaction then fails, the message is lost without a retry. The correct priority ordering for the consumer to be both durable and idempotent requires the transaction to commit before the acknowledgment is sent, and the idempotency key to be stable across all consumer invocations. Getting one right without the other produces a different failure mode.
Fixing Mode 1: derive the idempotency key from the message payload, not from UUID.randomUUID()
The fix is to replace UUID.randomUUID() with a content-hash key derived from the fields that uniquely identify the billing intent. Every redelivery of the same Kafka message carries the same payload; the content-hash produces the same key every time. The Stripe call with the same key on attempt 2 returns the same response as attempt 1 (the existing charge) without creating a new resource.
// BillingEventConsumer.java — fixed: content-hash idempotency key from message payload
@ApplicationScoped
public class BillingEventConsumer {
@Inject StripeClient stripeClient;
@Inject BillingAttemptRepository billingAttemptRepository;
@Incoming("billing-events")
@Transactional
public void onBillingEvent(BillingEvent event) {
// Content-hash key: same event fields always produce the same key.
// Redeliveries of the same Kafka message produce the same event fields.
// Stripe returns the existing charge on re-delivery — no ch_B.
String idempotencyKey = computeIdempotencyKey(
event.getCustomerId(),
event.getBillingPeriod(),
event.getAmountCents());
// Guard: if a BillingAttempt already exists with this key and SUCCEEDED status,
// skip the Stripe call entirely. This protects against redeliveries that arrive
// after a successful attempt was committed to the DB.
if (billingAttemptRepository.existsByIdempotencyKey(idempotencyKey)) {
return; // idempotent — already processed
}
ChargeParams params = ChargeParams.builder()
.customerId(event.getStripeCustomerId())
.amountCents(event.getAmountCents())
.currency("usd")
.idempotencyKey(idempotencyKey)
.build();
Charge charge = stripeClient.createCharge(params);
BillingAttempt attempt = new BillingAttempt(
event.getCustomerId(),
charge.getId(),
event.getBillingPeriod(),
idempotencyKey,
BillingStatus.SUCCEEDED);
billingAttemptRepository.persist(attempt);
}
private String computeIdempotencyKey(String customerId, String billingPeriod, long amountCents) {
// SHA-256 of the canonical billing inputs — deterministic and unique per billing intent.
String input = customerId + "|" + billingPeriod + "|" + amountCents;
try {
MessageDigest digest = MessageDigest.getInstance("SHA-256");
byte[] hash = digest.digest(input.getBytes(StandardCharsets.UTF_8));
return HexFormat.of().formatHex(hash).substring(0, 36);
} catch (NoSuchAlgorithmException e) {
throw new IllegalStateException("SHA-256 not available", e);
}
}
}
The DB-level guard (existsByIdempotencyKey) is an additional layer of defense. On a redelivery where the first attempt succeeded (Stripe committed ch_A and the database write also succeeded before the acknowledgment was lost), the guard short-circuits before calling Stripe again. This handles the acknowledgment-loss case: Kafka re-delivers the message even though the first processing was complete, because the acknowledgment (offset commit) was not received by the broker before the consumer crashed. Without the DB guard, even a content-hash key would make a redundant Stripe call (which Stripe would safely deduplicate), but the redundant Stripe HTTP request is avoidable overhead.
The order of operations also matters in the fixed version. The DB guard runs inside the @Transactional boundary before any Stripe call. If two concurrent redeliveries arrive simultaneously (possible in a partition-rebalanced Kafka setup), the unique constraint on idempotency_key in the database ensures exactly one of them commits the BillingAttempt row; the other sees a constraint violation and nacks, but because the key is content-based, the re-delivery would again find the row already present and return early.
Mode 2: @Incoming with Uni<Void> return type — UUID generated inside a Mutiny chain() callback is deferred until Uni subscription, so SmallRye re-subscription on nack produces UUID_B
Quarkus SmallRye Reactive Messaging supports reactive consumer methods that return Uni<Void>. When the consumer returns a Uni<Void>, SmallRye subscribes to the Uni and waits for its completion. If the Uni completes successfully (emits an item), SmallRye acks the message. If the Uni fails with an error, SmallRye nacks the message and, depending on the connector configuration, re-delivers it.
The failure mode in this section is not caused by message redelivery alone. It is caused by where the UUID is generated relative to the Mutiny evaluation model. Mutiny uses a deferred (lazy) execution model: operators like chain(), onItem().transformToUni(), and onItem().invoke() do not execute when the method body assembles the Uni chain. They execute when a subscriber subscribes to the Uni. For SmallRye consumers, subscription happens after the consumer method returns the Uni.
The developer writes:
// BillingEventConsumer.java — unsafe Mode 2: UUID inside Mutiny chain() callback
// UUID.randomUUID() is inside a lambda passed to chain().
// Lambdas passed to chain() are deferred suppliers: they execute at subscription time.
// SmallRye subscribes when the consumer method returns.
// On nack + redelivery: consumer method called again → new Uni → new subscription →
// chain() lambda re-executes → UUID.randomUUID() → UUID_B.
@ApplicationScoped
public class BillingEventConsumer {
@Inject StripeServiceClient stripeServiceClient;
@Inject BillingAttemptRepository billingAttemptRepository;
// Developer's reasoning:
// "UUID.randomUUID() is evaluated when onBillingEvent() is called.
// The lambda captures the UUID value — on re-delivery the method is called again
// but the UUID is still evaluated once per method call."
//
// The problem:
// UUID.randomUUID() is INSIDE the chain() lambda, not in the method body.
// The lambda is a deferred computation: it runs when SmallRye subscribes to the Uni.
// On re-delivery, SmallRye calls onBillingEvent() again → new Uni returned →
// SmallRye subscribes → chain() lambda executes → UUID.randomUUID() → UUID_B.
// This is equivalent to calling UUID.randomUUID() at method entry — UUID is
// re-generated per delivery attempt.
@Incoming("billing-events")
public Uni<Void> onBillingEvent(BillingEvent event) {
return Uni.createFrom().voidItem()
.chain(() -> {
// UUID generated here — inside a chain() lambda.
// This lambda is deferred: executed at Uni subscription time (not method call time).
// SmallRye subscribes AFTER this method returns.
// On re-delivery: new Uni → new subscription → this lambda re-executes → UUID_B.
String idempotencyKey = UUID.randomUUID().toString();
return stripeServiceClient.createCharge(
event.getStripeCustomerId(),
event.getAmountCents(),
idempotencyKey);
})
.chain(chargeId -> billingAttemptRepository.persistAsync(
event.getCustomerId(), chargeId, event.getBillingPeriod()))
.replaceWithVoid();
}
}
The developer’s mental model is that the UUID is generated “once per method call.” In imperative code, calling a method evaluates every statement in the method body immediately. In Mutiny reactive code, calling a method that returns Uni<Void> does not evaluate the operators: it assembles a pipeline description. The operators — chain(), onItem().invoke(), onFailure().recoverWithUni() — are functions that will be applied when the Uni is subscribed and items flow through the pipeline.
The UUID inside the chain() lambda is not different from UUID at method entry in terms of evaluation timing: both are evaluated exactly once per delivery attempt. But they feel different to the developer. A UUID at method entry is immediately visible as “evaluated when the method is called.” A UUID inside a lambda passed to chain() looks like it should be “captured” by the lambda closure. But there is nothing to capture: UUID.randomUUID() is not a variable being captured; it is a method call that is deferred into the lambda body and executed fresh every time the lambda is invoked.
If the developer had written:
// Contrast: UUID computed BEFORE the chain — still wrong, but for a different reason.
@Incoming("billing-events")
public Uni<Void> onBillingEvent(BillingEvent event) {
// UUID generated here — in the method body, before any chain() call.
// This IS evaluated at method call time, not at subscription time.
// BUT: on re-delivery, the consumer method is called again from line 1.
// UUID.randomUUID() still generates UUID_B on the second method call.
String idempotencyKey = UUID.randomUUID().toString();
return stripeServiceClient.createCharge(
event.getStripeCustomerId(),
event.getAmountCents(),
idempotencyKey)
.chain(chargeId -> billingAttemptRepository.persistAsync(
event.getCustomerId(), chargeId, event.getBillingPeriod()))
.replaceWithVoid();
}
This version evaluates the UUID at method-call time rather than at subscription time. It is still wrong: on redelivery, the consumer method is called again, and UUID.randomUUID() at method entry (or in the method body before the chain) still generates UUID_B. The distinction between “inside the chain lambda” and “outside the chain lambda in the method body” does not matter for correctness; both generate a fresh UUID per delivery. What matters is that the UUID must be computed from the message content, not from UUID.randomUUID(), so that all deliveries of the same message produce the same key.
The subtle Mutiny-specific variant that is worth highlighting separately involves Uni.createFrom().item(supplier). A developer might write:
// Variant: UUID inside Uni.createFrom().item(Supplier) — deferred evaluation
// This is MORE deferred than UUID in the method body: evaluated at subscription time.
// The developer thinks this is "lazy evaluation of the billing request" but it is
// lazy UUID generation — re-evaluated per subscription, per delivery.
@Incoming("billing-events")
public Uni<Void> onBillingEvent(BillingEvent event) {
return Uni.createFrom().item(() -> {
// Supplier evaluated at subscription time.
// Each SmallRye subscription (one per delivery) evaluates this.
String idempotencyKey = UUID.randomUUID().toString();
return new ChargeRequest(event.getStripeCustomerId(),
event.getAmountCents(), idempotencyKey);
})
.chain(req -> stripeServiceClient.createChargeAsync(req))
.chain(chargeId -> billingAttemptRepository.persistAsync(
event.getCustomerId(), chargeId, event.getBillingPeriod()))
.replaceWithVoid();
}
In this variant the UUID is even more explicitly deferred: it is inside a Supplier passed to Uni.createFrom().item(Supplier). The supplier is evaluated when the Uni is subscribed, not when it is constructed. SmallRye subscribes when the consumer method returns the Uni. On redelivery, the consumer method is called again, returns a new Uni.createFrom().item(supplier), SmallRye subscribes, the supplier evaluates, UUID_B is generated.
All three variants — UUID in a chain() lambda, UUID in the method body, UUID in a Uni.createFrom().item(Supplier) — produce the same result: UUID_B on any delivery after the first. The correct solution is identical for all variants.
Fixing Mode 2: compute the content-hash key before the Uni chain, capture it as a final local variable
The fix is structurally simple: compute the content-hash idempotency key before the Uni pipeline is assembled, assign it to a final local variable, and reference the variable inside the chain. The key is now stable: it is computed once per message payload (not per subscription), and all redeliveries of the same Kafka message produce the same payload and therefore the same key.
// BillingEventConsumer.java — fixed Mode 2: content-hash key computed before Uni chain
@ApplicationScoped
public class BillingEventConsumer {
@Inject StripeServiceClient stripeServiceClient;
@Inject BillingAttemptRepository billingAttemptRepository;
@Incoming("billing-events")
public Uni<Void> onBillingEvent(BillingEvent event) {
// Content-hash key computed BEFORE the chain — eagerly, at method call time.
// Captured as final — the same value is referenced by all chain lambdas.
// On re-delivery: same event payload → same key computation → same idempotencyKey value.
// Stripe returns the existing charge on any re-delivery after the first successful call.
final String idempotencyKey = computeIdempotencyKey(
event.getCustomerId(),
event.getBillingPeriod(),
event.getAmountCents());
return billingAttemptRepository.existsByIdempotencyKeyAsync(idempotencyKey)
.chain(alreadyProcessed -> {
if (alreadyProcessed) {
// Idempotent: this delivery was already processed successfully.
return Uni.createFrom().voidItem();
}
// idempotencyKey captured from outer scope — stable across all subscriptions.
return stripeServiceClient.createCharge(
event.getStripeCustomerId(),
event.getAmountCents(),
idempotencyKey)
.chain(chargeId -> billingAttemptRepository.persistAsync(
event.getCustomerId(), chargeId,
event.getBillingPeriod(), idempotencyKey));
})
.replaceWithVoid();
}
private String computeIdempotencyKey(String customerId, String billingPeriod, long amountCents) {
String input = customerId + "|" + billingPeriod + "|" + amountCents;
try {
MessageDigest digest = MessageDigest.getInstance("SHA-256");
byte[] hash = digest.digest(input.getBytes(StandardCharsets.UTF_8));
return HexFormat.of().formatHex(hash).substring(0, 36);
} catch (NoSuchAlgorithmException e) {
throw new IllegalStateException("SHA-256 not available", e);
}
}
}
The key point in this fixed version: idempotencyKey is declared final and computed in the method body before the return statement. It is captured by the chain() lambda as a closed-over variable. Lambda closures in Java capture the value at the time the lambda is created (which is at method execution time, when the lambda expression is evaluated). The value captured is the content-hash string, not the expression UUID.randomUUID(). On re-delivery, the consumer method is called again, the content-hash produces the same string (same message payload), and the same string is captured by the new lambda. The Stripe call with the same key returns the existing charge or is deduplicated by the database guard before the Stripe call is made.
There is also a Quarkus Panache Reactive variant worth noting: when using @WithTransaction (Quarkus’s reactive transaction annotation) on a consumer that returns Uni<Void>, the transaction wraps the entire Uni pipeline. The UUID must still be computed outside the chain to be stable. @WithTransaction does not change the deferred-evaluation semantics of Mutiny operators; it only adds transaction begin/commit/rollback around the Uni subscription lifecycle.
Mode 3: @Outgoing Kafka emitter inside @Transactional — Kafka messages are committed to the broker before the application transaction commits or rolls back, producing phantom messages that drive duplicate consumer executions
The third failure mode does not involve a consumer retry at all. It involves the write ordering between a Kafka producer and a JPA transaction manager. The developer writes a Quarkus service method that performs two writes: a JPA database write (persisting a BillingRecord entity) and a Kafka message emission (sending a BillingEvent to the billing-events topic). Both writes are inside a @Transactional method. The developer expects both writes to be atomic: either both commit or both are rolled back.
This expectation is correct when the two resources are enrolled in the same XA distributed transaction managed by a JTA transaction manager. JPA, backed by a JDBC datasource that supports XA, can participate in JTA. However, Quarkus’s default Kafka client (backed by the Confluent Kafka producer) is not an XA resource and is not enrolled in the JTA transaction by default. The Kafka producer sends messages to the broker independently of the JPA transaction’s two-phase commit. When the developer calls emitter.send(event), SmallRye Reactive Messaging routes the message to the Kafka producer, which buffers it and sends it to the broker. The JPA transaction has not committed yet; the Kafka message may be sent to the broker before or concurrently with the JPA commit.
// BillingService.java — unsafe Mode 3: @Outgoing Kafka emitter inside @Transactional
// Kafka producer is not a JTA resource.
// emitter.send() commits the message to the Kafka broker independently of the @Transactional.
// If @Transactional later rolls back (validation failure, constraint violation, etc.),
// the Kafka message is already in the broker.
// @Incoming consumer receives the phantom message → UUID → Stripe → ch_A.
// Application retry (unaware of the phantom message) emits again → ch_B.
@ApplicationScoped
public class BillingService {
@Inject BillingRecordRepository billingRecordRepository;
// Channel "billing-events-out" mapped to Kafka topic "billing-events" in application.properties.
@Channel("billing-events-out")
MutinyEmitter<BillingEvent> billingEventEmitter;
// Developer's reasoning:
// "@Transactional makes the DB write and the Kafka emit atomic.
// If the validation fails after emitter.send(), the @Transactional rolls back both.
// The Kafka message is not committed to the broker if the transaction rolls back."
//
// The problem:
// Kafka producer is not enrolled in JTA.
// emitter.send(event) sends the message to the Kafka broker independently.
// The message is in the broker (on the topic partition) whether or not the JPA
// transaction commits or rolls back.
// If postValidate() throws, @Transactional rolls back the BillingRecord insert.
// But the "billing-events" Kafka message is already committed to the broker.
// @Incoming BillingEventConsumer receives the message → UUID → Stripe → ch_A.
// Application caller catches the exception and retries initiateBillingRun().
// initiateBillingRun() calls emitter.send(event) again → second Kafka message.
// Consumer processes second message → UUID_B → Stripe → ch_B.
@Transactional
public void initiateBillingRun(String customerId, long amountCents, String billingPeriod) {
// Step 1: persist the BillingRecord.
BillingRecord record = new BillingRecord(customerId, amountCents, billingPeriod,
BillingStatus.PENDING);
billingRecordRepository.persist(record);
// Step 2: emit the Kafka message.
// Developer expects this to be atomic with the DB write above.
// In reality: Kafka producer sends to broker now, independently of JTA transaction.
BillingEvent event = new BillingEvent(record.getId(), customerId,
record.getStripeCustomerId(), amountCents, billingPeriod);
billingEventEmitter.sendAndAwait(event);
// Step 3: post-validation that may throw.
// If this throws: @Transactional rolls back the BillingRecord insert.
// The Kafka message is already in the broker — NOT rolled back.
postValidate(customerId, billingPeriod);
}
private void postValidate(String customerId, String billingPeriod) {
// Example: check that billing hasn't already been completed for this period.
// Could throw IllegalStateException if billing already run this period.
// ...
}
}
The developer who does not know about Kafka’s transactional producer API may assume that @Transactional rolls back the Kafka emission. This assumption is natural: @Transactional is a Java EE / Jakarta EE construct, and in a full Jakarta EE application server with XA datasources, you can enroll message producers in the JTA transaction (JMS with XA is a classic example). Kafka, however, is not JMS. The SmallRye Reactive Messaging Kafka connector uses the standard Confluent Kafka producer, which supports its own transaction protocol (Kafka’s exactly-once semantics using transactional.id) but does not participate in an external JTA transaction unless you configure the smallrye.kafka.transactions.enabled property and use Quarkus’s Kafka Transactions extension.
In the default Quarkus configuration, when the service method calls emitter.sendAndAwait(event), SmallRye buffers and then flushes the Kafka message. The Kafka producer commits the message to the partition. The offset for the message is assigned. Any consumer polling that partition can receive the message. The JPA transaction for the current thread has not yet committed: the BillingRecord row is not yet visible to other transactions. But the Kafka consumer downstream can receive and begin processing the message before the producer’s JPA transaction commits or rolls back.
When postValidate() throws, the @Transactional interceptor rolls back the JPA transaction. The BillingRecord insert is undone. But the Kafka message is in the broker and the consumer has likely already received it. The consumer generates UUID.randomUUID() (the unsafe pattern from Mode 1), calls Stripe, and creates ch_A. The service caller catches the exception, retries initiateBillingRun() (perhaps via a user-facing API retry or a scheduled job that retries failed billing runs). The second call to initiateBillingRun() persists a new BillingRecord and emits another Kafka message. The consumer processes the second message with UUID_B and creates ch_B.
Even if the consumer in this scenario uses a content-hash idempotency key (the fix from Mode 1), Mode 3 produces a different kind of double-processing. The first Kafka message was produced before the JPA transaction committed; the BillingRecord row was rolled back. The consumer processes the message, persists a BillingAttempt (the consumer has its own @Transactional), and calls Stripe. The content-hash idempotency key guards against a second Stripe call for the same event. But the first BillingRecord (from the producer’s rolled-back transaction) never exists; the BillingAttempt created by the consumer references a BillingRecord that does not exist. The data is inconsistent: Stripe has a charge for a billing run that has no corresponding BillingRecord in the database.
The second retry call to initiateBillingRun() persists a new BillingRecord with a new UUID primary key and emits a second Kafka message with a different record.getId(). If the consumer’s content-hash key is derived from customerId + billingPeriod + amountCents (not from record.getId()), the content-hash key is the same as the first message. The consumer finds the existing BillingAttempt and skips the Stripe call. One charge, but the data model has a dangling BillingAttempt linked to a non-existent BillingRecord.
If the content-hash key includes record.getId(), the second message’s key is different (new record ID). The consumer processes both messages, and depending on timing, two BillingAttempt rows may be created with different idempotency keys but the same Stripe customer, amount, and billing period. Stripe deduplication is by Idempotency-Key header value; the second call with a different key creates ch_B.
Fixing Mode 3: emit Kafka messages inside Kafka’s own transaction, or use transactional outbox
There are two patterns to fix Mode 3.
Option A: Kafka Transactions (Kafka’s exactly-once semantics). Configure the Quarkus Kafka Transactions extension. This enrolls the Kafka producer in a Kafka transaction (separate from JTA) that is committed and rolled back coordinately with the application’s JTA transaction via a Quarkus interceptor that registers a JTA synchronization. When the JTA transaction rolls back, the Kafka transaction is also aborted, and the broker discards the message. Consumers configured with isolation.level=read_committed do not see messages from aborted Kafka transactions.
# application.properties — Kafka Transactions configuration for Mode 3 fix
# Requires: io.quarkus:quarkus-smallrye-reactive-messaging-kafka-transactions
# Enable Kafka transaction support for the outgoing channel.
mp.messaging.outgoing.billing-events-out.transactional=true
mp.messaging.outgoing.billing-events-out.transactional-id=billing-events-producer
# Consumer: read_committed isolation — do not consume from aborted Kafka transactions.
mp.messaging.incoming.billing-events.isolation.level=read_committed
# Producer: all acknowledgments required for exactly-once.
mp.messaging.outgoing.billing-events-out.acks=all
// BillingService.java — Kafka Transaction fix for Mode 3
// @KafkaTransaction registers a JTA synchronization:
// JTA commit → Kafka transaction commit; JTA rollback → Kafka transaction abort.
// Consumer with read_committed sees only committed Kafka messages.
@ApplicationScoped
public class BillingService {
@Inject BillingRecordRepository billingRecordRepository;
@Channel("billing-events-out")
@KafkaClientService // or use KafkaTransactions injection directly
MutinyEmitter<BillingEvent> billingEventEmitter;
@Transactional
public void initiateBillingRun(String customerId, long amountCents, String billingPeriod) {
BillingRecord record = new BillingRecord(customerId, amountCents, billingPeriod,
BillingStatus.PENDING);
billingRecordRepository.persist(record);
// With Kafka Transactions enabled and isolation.level=read_committed on the consumer:
// - If postValidate() throws and @Transactional rolls back the JPA transaction,
// the JTA synchronization also aborts the Kafka transaction.
// - The broker discards the message for this producer transaction.
// - Consumer never receives it.
BillingEvent event = new BillingEvent(record.getId(), customerId,
record.getStripeCustomerId(), amountCents, billingPeriod);
billingEventEmitter.sendAndAwait(event);
postValidate(customerId, billingPeriod);
}
}
Option B: Transactional outbox pattern. Instead of emitting the Kafka message directly inside the service method, write an OutboxEvent row to the database inside the same JPA transaction as the BillingRecord. A separate process (Debezium CDC, or a Quarkus @Scheduled poller) reads committed OutboxEvent rows and publishes them to Kafka. Because the OutboxEvent is written in the same JPA transaction as the BillingRecord, they are either both committed or both rolled back. The Kafka message is only produced after the database row is committed; there are no phantom messages.
// BillingService.java — Transactional Outbox fix for Mode 3
// OutboxEvent is written in the same @Transactional as BillingRecord.
// If postValidate() throws: both BillingRecord and OutboxEvent are rolled back.
// No phantom Kafka message. No phantom consumer invocation. No ch_A from rolled-back data.
@ApplicationScoped
public class BillingService {
@Inject BillingRecordRepository billingRecordRepository;
@Inject OutboxEventRepository outboxEventRepository;
@Transactional
public void initiateBillingRun(String customerId, long amountCents, String billingPeriod) {
BillingRecord record = new BillingRecord(customerId, amountCents, billingPeriod,
BillingStatus.PENDING);
billingRecordRepository.persist(record);
// OutboxEvent written in the same transaction as BillingRecord.
// Both committed atomically, or both rolled back.
// No Kafka message emitted here — the outbox poller/CDC handles that.
OutboxEvent outboxEvent = new OutboxEvent(
"billing-events",
record.getId().toString(),
serializeToJson(new BillingEvent(record.getId(), customerId,
record.getStripeCustomerId(), amountCents, billingPeriod)));
outboxEventRepository.persist(outboxEvent);
// If postValidate() throws: @Transactional rolls back BillingRecord AND OutboxEvent.
// Outbox poller never sees the OutboxEvent row. No Kafka message. No consumer. No Stripe call.
postValidate(customerId, billingPeriod);
}
}
The transactional outbox pattern requires more infrastructure (a Debezium connector or a polling loop) but works with any JTA-capable database without requiring Kafka’s exactly-once semantics configuration. For teams that are not ready to configure Kafka transactions, the outbox is the more universally available option. For teams already operating Kafka with exactly-once semantics in their pipeline (to handle consumer exactly-once processing), enabling the Quarkus Kafka Transactions integration is the lower-infrastructure path.
Both options are compatible with the content-hash idempotency key fix from Mode 1. Even after fixing Mode 3 so that phantom Kafka messages are never produced, the consumer should still use a content-hash key. Mode 3 produces duplicate consumer invocations from phantom messages; Mode 1 produces duplicate consumer invocations from message redeliveries on nack. Both scenarios lead to duplicate Stripe calls unless the consumer’s idempotency key is stable across invocations.
Comparison table: three SmallRye Reactive Messaging failure modes
| Mode | Root cause | Trigger | Why developer misses it | Detection signal | Fix |
|---|---|---|---|---|---|
1. @Incoming + @Transactional |
Stripe committed before DB write; @Transactional rolls back DB but not Stripe; SmallRye nacks; redelivery generates UUID_B |
DB write failure (constraint, optimistic lock, connection loss) after Stripe call succeeds | “@Transactional on the consumer makes processing atomic — rollback means the retry starts clean including Stripe” |
Stripe Dashboard: two PaymentIntents for same customer + billing period within seconds; DB: BillingAttempt with UUID_B but no UUID_A record |
Content-hash key from message payload; DB guard existsByIdempotencyKey before Stripe call |
2. Uni<Void> + Mutiny deferred UUID |
UUID inside chain() lambda or Uni.createFrom().item(Supplier) is deferred to subscription time; each SmallRye redelivery creates a new Uni and re-subscribes |
Failed Uni (Stripe error, DB error, timeout) → SmallRye nacks → redelivery → new Uni → new subscription → chain() re-executes → UUID_B | “UUID is generated when the consumer method is called — lambda captures the value; on re-delivery the same captured value is used” | Stripe Dashboard: duplicate PaymentIntents within delivery retry window; application log: multiple Idempotency-Key values for same Kafka offset partition key |
Compute content-hash key before chain assembly; assign to final local variable; lambda captures the string value, not the expression |
3. @Outgoing + @Transactional without Kafka TX |
Kafka producer not enrolled in JTA; emitter.send() commits message to broker before JPA transaction commits or rolls back; phantom messages reach consumer after JPA rollback |
JPA transaction rollback after emitter.send() — validation failure, constraint violation, etc. |
“@Transactional makes the Kafka emit and DB write atomic — if the transaction rolls back, the Kafka message is discarded too” |
Consumer log: processing message for billing period that has no corresponding BillingRecord; Stripe charge exists but no BillingRecord row in DB; consumer offset ahead of producer transaction outcome |
Enable Kafka Transactions extension (Quarkus) with transactional=true + read_committed consumer; or use Transactional Outbox pattern |
Testing: how to write QuarkusTest consumer tests that catch all three modes
Each of the three failure modes requires a different testing approach because each is triggered by a different mechanism. Mode 1 requires a test that simulates a DB write failure after a successful Stripe call and verifies the consumer does not double-charge on redelivery. Mode 2 requires a test that verifies the idempotency key is stable across multiple SmallRye subscriptions (simulating redelivery). Mode 3 requires a test that verifies the Kafka message is not produced if the JPA transaction rolls back.
// BillingEventConsumerTest.java — QuarkusTest covering Modes 1 and 2
@QuarkusTest
@QuarkusTestResource(WireMockTestResource.class) // Stripe mock
class BillingEventConsumerTest {
// Inject the in-memory channel emitter to send test messages directly.
@Inject
@Channel("billing-events")
Emitter<BillingEvent> testEmitter;
@Inject BillingAttemptRepository billingAttemptRepository;
@InjectMock StripeClient stripeClient;
@BeforeEach
@Transactional
void cleanDb() {
billingAttemptRepository.deleteAll();
}
@Test
void mode1_redeliveryDoesNotDoubleCharge() throws InterruptedException {
// Arrange: first Stripe call succeeds; DB persist will fail on first delivery
// (simulated by a duplicate key already in DB for this idempotency key).
// The second delivery (redelivery after nack) should find the key in DB
// and skip the Stripe call.
// Insert a pre-existing BillingAttempt with the content-hash key
// to simulate a redelivery after a successful first attempt.
String expectedKey = computeContentHashKey("cust_001", "2026-10", 1000L);
// Pre-insert the idempotency key record to simulate already-processed state.
insertBillingAttempt("cust_001", "2026-10", "ch_existing_001", expectedKey);
// Act: send the message (simulates redelivery of an already-processed message).
BillingEvent event = new BillingEvent("cust_001", "stripe_cus_001",
"2026-10", 1000L);
testEmitter.send(event).await().atMost(Duration.ofSeconds(5));
// Allow consumer to process.
Thread.sleep(200);
// Assert: Stripe was NOT called (consumer found existing key and returned early).
verify(stripeClient, never()).createCharge(any());
// Assert: still exactly one BillingAttempt row (no duplicate inserted).
assertThat(billingAttemptRepository.count()).isEqualTo(1L);
}
@Test
void mode2_idempotencyKeyStableAcrossRedeliveries() throws Exception {
// Arrange: Stripe succeeds; capture all calls to verify key stability.
List<String> capturedKeys = new CopyOnWriteArrayList<>();
when(stripeClient.createCharge(any())).thenAnswer(invocation -> {
ChargeParams params = invocation.getArgument(0);
capturedKeys.add(params.getIdempotencyKey());
return new Charge("ch_test_" + capturedKeys.size());
});
// Simulate two deliveries of the same logical billing event.
BillingEvent event = new BillingEvent("cust_002", "stripe_cus_002", "2026-10", 2000L);
// First delivery.
testEmitter.send(event).await().atMost(Duration.ofSeconds(5));
Thread.sleep(200);
// Second delivery (same payload — simulates Kafka redelivery).
testEmitter.send(event).await().atMost(Duration.ofSeconds(5));
Thread.sleep(200);
// With content-hash key: both deliveries produce the same key.
// DB guard on second delivery: key already exists, Stripe not called again.
// Stripe called exactly once.
assertThat(capturedKeys).hasSize(1);
// With UUID.randomUUID(): capturedKeys would have size 2 with different values.
// This assertion fails and exposes the bug.
if (capturedKeys.size() > 1) {
assertThat(capturedKeys.get(0))
.as("Idempotency-Key must be identical across all deliveries of the same event")
.isEqualTo(capturedKeys.get(1));
}
}
}
// BillingServiceKafkaTransactionTest.java — QuarkusTest covering Mode 3
@QuarkusTest
class BillingServiceKafkaTransactionTest {
@Inject BillingService billingService;
@Inject BillingRecordRepository billingRecordRepository;
// Capture outbox events written to DB instead of testing Kafka directly.
// With the Transactional Outbox pattern, a rolled-back @Transactional
// means no OutboxEvent row — which means no Kafka message — which means no consumer.
@Inject OutboxEventRepository outboxEventRepository;
@BeforeEach
@Transactional
void cleanDb() {
billingRecordRepository.deleteAll();
outboxEventRepository.deleteAll();
}
@Test
void mode3_kafkaMessageNotProducedWhenTransactionRollsBack() {
// Arrange: initiateBillingRun will call postValidate which throws.
// With the Transactional Outbox fix: neither BillingRecord nor OutboxEvent committed.
assertThatThrownBy(() ->
billingService.initiateBillingRun("cust_003", 3000L, "already-billed-period"))
.isInstanceOf(IllegalStateException.class);
// Assert: no BillingRecord committed.
assertThat(billingRecordRepository.count()).isEqualTo(0L);
// Assert: no OutboxEvent committed (no phantom Kafka message).
// If this assertion fails with count = 1, the outbox write was committed
// independently of the @Transactional rollback — phantom message produced.
assertThat(outboxEventRepository.count())
.as("OutboxEvent must be rolled back with the @Transactional — no phantom Kafka message")
.isEqualTo(0L);
}
@Test
void mode3_successfulTransactionProducesExactlyOneOutboxEvent() {
// Arrange: valid billing run — no validation failure.
billingService.initiateBillingRun("cust_004", 4000L, "2026-10");
// Assert: exactly one BillingRecord and exactly one OutboxEvent.
assertThat(billingRecordRepository.count()).isEqualTo(1L);
assertThat(outboxEventRepository.count()).isEqualTo(1L);
// The OutboxEvent contains the billing event payload that the outbox poller
// will eventually publish to Kafka — after the JPA transaction commits.
OutboxEvent outboxEvent = outboxEventRepository.findAll().firstResult();
assertThat(outboxEvent.getAggregateId())
.isEqualTo(billingRecordRepository.findAll().firstResult().getId().toString());
}
}
These tests use @QuarkusTest with an in-memory SmallRye channel (smallrye-reactive-messaging-in-memory connector) for Modes 1 and 2, avoiding the need for a running Kafka broker in the test environment. The in-memory channel delivers messages synchronously or near-synchronously, making the timing assertions in the tests reliable without extensive polling.
For Mode 3, the test validates the transactional boundary at the database level (OutboxEvent count) rather than at the Kafka level. This is more reliable in a test environment: asserting that a Kafka message was or was not produced requires either a real Kafka broker (via Testcontainers) or a SmallRye in-memory channel spy. Asserting that an OutboxEvent row was or was not committed is a straightforward database assertion that works with the standard Quarkus test datasource.
Required assertions checklist for any Quarkus SmallRye Reactive Messaging Stripe integration test
- Idempotency key stability across redeliveries. Send the same message to the in-memory channel twice. Assert that Stripe was called at most once (the second delivery found the key in the DB and returned early). If Stripe was called twice, assert
capturedKeys.get(0).equals(capturedKeys.get(1)). - Content-hash key derivation test. Assert that
computeIdempotencyKey(customerId, billingPeriod, amountCents)returns the same value when called twice with the same inputs. Assert that it returns a different value when called with a differentbillingPeriod. - DB guard idempotency test. Pre-insert a
BillingAttemptwith the expected content-hash key. Send the message. Assert Stripe was not called. Assert DB still has exactly oneBillingAttemptrow. - Transaction rollback produces no OutboxEvent (Mode 3 outbox). Call the service method knowing validation will throw. Assert
outboxEventRepository.count() == 0. If using Kafka Transactions instead of outbox: assert Stripe consumer receives zero messages after the rollback. - Kafka isolation level test (Mode 3 Kafka TX). If using Kafka Transactions: configure a test consumer with
isolation.level=read_committed. Produce a message inside a Kafka transaction, abort the transaction, poll the consumer. Assert the consumer received zero messages (aborted transaction messages are not returned toread_committedconsumers).
Detection signals in production: identifying SmallRye Reactive Messaging Stripe duplicates
The three failure modes produce overlapping but distinct observable signals in production. Knowing which signal maps to which mode helps narrow the investigation quickly.
Mode 1 signals (nack + redelivery, UUID_B from method entry):
- Stripe Dashboard shows two
PaymentIntentorChargeobjects for the same customer and billing period, created within seconds of each other. The time delta between ch_A and ch_B is the sum of the message redelivery backoff and the consumer processing time. - Application DB: a
BillingAttemptrow with statusSUCCEEDEDandidempotency_key = UUID_B, but no correspondingBillingAttemptrow withUUID_A. The UUID_A attempt was rolled back with the first failed transaction. - Consumer logs:
onBillingEvent called for customer=cust_001 billing_period=2026-10appearing twice in quick succession on the same Kafka partition key, with differentidempotencyKeylog fields.
Mode 2 signals (Mutiny deferred UUID, subscription-time evaluation):
- Stripe duplicate signals are the same as Mode 1: two charges within seconds. The distinguishing signal is in the application log: if the consumer logs the idempotency key at the point it is computed, Mode 2 produces two log lines with different keys from the same consumer method at nearly the same timestamp.
- If the consumer logs the Uni pipeline evaluation trace (using
.invoke(() -> log.info("evaluating chain()"))), Mode 2 shows the chain evaluation happening twice: once per subscription (once per SmallRye delivery attempt).
Mode 3 signals (phantom Kafka message from rolled-back @Transactional):
- Consumer receives and processes a message for a
BillingRecordthat does not exist in the database at the time of consumer processing. Consumer logs:BillingRecordNotFoundException for record_id=.... If the consumer has afindByIdguard, it logs “billing record not found” and nacks the message, leading to Mode 1-style redelivery loops. - A Stripe charge exists with a Stripe customer ID and amount that does not correspond to any committed
BillingRecordrow. This is auditable by joining Stripe Dashboard export data with the production DB. - If the transactional outbox is in use and you observe OutboxEvent rows without corresponding
BillingRecordrows, the outbox write and the BillingRecord write are not in the same transaction — the outbox pattern is implemented incorrectly.
All three modes are auditable by Keybrake’s Stripe proxy audit log: the proxy records every request forwarded to Stripe, including the Idempotency-Key header value, the Stripe customer ID, the amount, the response status, and the timestamp. Querying the audit log for duplicate (customer_id, amount, billing_period) tuples with different idempotency_key values identifies the production failure in the same pass, regardless of which mode caused it. The Keybrake spend cap can be configured to alert when a second charge attempt for the same customer + billing period arrives within 10 seconds — which is the typical window for all three SmallRye reactive messaging failure modes above.
Stop duplicate Stripe charges before they hit your customers
Keybrake proxies your agent’s Stripe calls and enforces per-vendor spend caps, allowlists, and an audit log. One-click revoke if a consumer loop runs away.