Project Reactor and Kotlin Flow Stripe Integration: How Mono.defer{} Inside retryWhen() Re-evaluates UUID on Re-subscription, Kotlin Flow Cold Builders Re-execute UUID Generation on Flow.retryWhen{} Re-collection, and TransactionalOperator vs. @Transactional suspend fun Both Lose the Idempotency Key When @Retryable Is the Outer AOP Proxy
Project Reactor and Kotlin Flow are the two dominant reactive abstractions for Spring Boot and Kotlin backends. They share a root problem for Stripe idempotency: any re-subscription or re-collection caused by a retry operator re-executes the cold source, including any UUID.randomUUID() call placed inside it. The failure modes differ in mechanism — Reactor’s publisher factory functions vs. Flow’s cold builder semantics vs. Spring’s AOP proxy layer — but the outcome is the same: UUID_B reaches Stripe on attempt 2, producing a duplicate charge alongside the already-committed ch_A.
This post covers three failure modes, one specific to Project Reactor, one specific to Kotlin Flow, and one cross-framework comparison that applies to both. It is structurally distinct from the Spring Data R2DBC + Kotlin Coroutines post (which covers @Transactional suspend fun + @Retryable fresh coroutine, flatMap placement of retryWhen(), and TransactionalOperator wrapping UUID generation in a single framework context) and from the Spring WebFlux Functional Endpoints post (which covers bodyToMono().cache(), WebClient header customizer, and Mono.fromCallable() in per-item retry). The modes here are specific to how Mono.defer{} interacts with Reactor’s re-subscription model, how Kotlin Flow’s cold builder semantics differ from Reactor’s cold publisher in batch billing scenarios, and how TransactionalOperator (Reactor) and @Transactional suspend fun (Kotlin) produce the same UUID re-generation problem through mechanically different paths when @Retryable is the outer AOP proxy.
Background: what “cold” means in Reactor vs. Kotlin Flow, and why it matters for Stripe
In Project Reactor, a Mono or Flux is a description of a computation, not an execution of it. Subscribing to a Mono or Flux executes the pipeline. A “cold” publisher runs its producer logic independently for each subscriber. Mono.fromCallable(), Mono.defer(), Flux.generate(), and most Reactor operators are cold by default — they execute their factory functions or lambdas per subscription. A “hot” publisher shares a single execution across all subscribers. Mono.just(value) captures a value eagerly at assembly time; subscribing to a Mono.just() does not re-evaluate the captured value, making it “hot” with respect to the value, though it is still cold with respect to delivering the value to each subscriber separately.
In Kotlin Flow, all flows built with the flow{}, channelFlow{}, or callbackFlow{} builders are cold by default: the builder block executes independently for each call to collect(). SharedFlow and StateFlow are hot. Cold flow semantics mean that calling retryWhen{}` on a cold flow re-collects from the upstream producer, re-executing all emit() calls in the builder block.
For Stripe idempotency, the critical consequence is: any UUID.randomUUID() call placed inside a cold publisher factory function (Reactor defer{}, fromCallable{}) or inside a cold flow builder block (flow{}, channelFlow{}) will be re-evaluated on every re-subscription or re-collection caused by a retry operator. The Stripe Idempotency-Key header sent in attempt 2 will be UUID_B, not UUID_A. Stripe sees two distinct keys for the same logical charge and commits two PaymentIntents.
Stripe’s idempotency contract: a POST to any Stripe mutating endpoint with an Idempotency-Key header deduplicates requests to that endpoint per key per API key for 24 hours. The key is the only deduplication mechanism. Stripe does not inspect customer ID, amount, or metadata to detect duplicates across different keys.
Mode 1 (Reactor): Mono.defer{} inside retryWhen() — defer re-evaluates UUID on every re-subscription
The developer is building a reactive billing service using Project Reactor. They know that Mono.just(UUID.randomUUID()) evaluates the UUID at assembly time — the moment the Mono.just() call executes in the method that builds the chain. They correctly reason that if the chain is assembled once and subscribed multiple times (e.g., via retryWhen()), Mono.just() captures the UUID from assembly time and that same UUID is delivered to every subscriber. This is true.
However, when the chain is assembled inside a method that is called per-request (a Spring WebFlux handler or a service method returning a Mono), the entire chain is assembled fresh on each incoming request. The developer recognizes this and switches to Mono.defer{}: they wrap the UUID generation in a defer block to make it “reactive-native” and avoid blocking assembly. The reasoning: Mono.defer{} defers evaluation to subscription time, which is the reactive idiom for lazy initialization.
The problem: defer defers evaluation to subscription time, and retryWhen() causes re-subscription. When the Stripe HTTP call fails and retryWhen() fires, it re-subscribes the entire upstream chain. The defer block re-executes its factory function. UUID.randomUUID() generates UUID_B. The flatMap downstream of the defer sends UUID_B to Stripe as the Idempotency-Key header. ch_B is created in Stripe alongside ch_A, which was already committed before the first attempt failed.
The developer misconception takes a specific form: “Mono.defer{} lazily wraps UUID generation — retryWhen() retries the operator that threw the exception, not the deferred source that had already been evaluated and delivered its value downstream.” The developer is conflating delivery of the value with evaluation of the factory. When retryWhen() re-subscribes, it re-subscribes from the beginning of the chain it wraps. If the defer is inside that chain, the defer’s factory is re-executed. The value is not cached between subscriptions.
// BillingService.java — unsafe Mode 1: Mono.defer{} inside retryWhen() chain
@Service
public class ReactiveBillingService {
private final R2dbcBillingRepository repository;
private final StripeWebClient stripeClient;
// Developer's reasoning:
// "Mono.just(UUID.randomUUID()) evaluates the UUID at assembly time, which is
// wrong for a reactive chain — the UUID should be evaluated at subscription time
// (per request). Mono.defer{} defers evaluation to subscription time. retryWhen()
// retries from the operator that threw (the Stripe HTTP call), not from my defer
// block that had already evaluated and passed the UUID downstream."
//
// The problem:
// retryWhen() re-subscribes the ENTIRE upstream operator chain it wraps.
// Mono.defer{} is part of that chain. On re-subscription, defer{}'s factory
// function re-executes. UUID.randomUUID() generates UUID_B on attempt 2.
//
// Timeline:
// Assembly: chain is built (no UUID generated yet — defer is lazy).
//
// Subscription 1 (attempt 1):
// defer{} factory executes → UUID_A = UUID.randomUUID()
// flatMap: BillingAttempt saved with UUID_A in R2DBC transaction T1.
// flatMap: Stripe called with Idempotency-Key: UUID_A → ch_A committed in Stripe.
// Stripe returns 500 (transient error).
// retryWhen() catches the error signal — decides to retry.
//
// Re-subscription (attempt 2):
// retryWhen() re-subscribes from the TOP of the chain it wraps.
// defer{} factory re-executes → UUID_B = UUID.randomUUID(). ← NEW UUID
// flatMap: new BillingAttempt INSERT with UUID_B (T1 may have rolled back; T2 opens).
// flatMap: Stripe called with Idempotency-Key: UUID_B → ch_B committed. ← DUPLICATE
//
// Result: ch_A and ch_B both committed. Customer charged twice.
public Mono chargeCustomer(String customerId, String stripeCustomerId,
long amountCents, String billingPeriod) {
return Mono.defer(() -> {
// Developer places UUID generation here for "reactive-correct lazy evaluation."
String idempotencyKey = UUID.randomUUID().toString();
return Mono.just(idempotencyKey);
})
.flatMap(key -> repository.save(
BillingAttempt.pending(customerId, amountCents, billingPeriod, key))
.thenReturn(key))
.flatMap(key -> stripeClient.createPaymentIntent(
stripeCustomerId, amountCents, key))
.retryWhen(Retry.backoff(3, Duration.ofMillis(200))
.filter(ex -> ex instanceof WebClientResponseException.InternalServerError));
}
}
Why Mono.defer{} does not protect UUID across retryWhen()
Mono.defer(supplier) creates a new publisher by calling the supplier on every subscription. The supplier is a factory: it runs once per subscription. When retryWhen() decides to retry, it signals the upstream to cancel and then re-subscribes. “Re-subscribe” means calling subscribe() on the upstream Mono again, which triggers the defer’s supplier again. This is by design — defer is documented as “Create a Mono provider that will supply a target Mono to subscribe to for each subscriber.” The phrase “for each subscriber” is the operative clause.
The distinction the developer misses is between operator re-subscription and partial chain re-execution from the point of failure. Reactor does not support the latter. When retryWhen() retries, it always re-subscribes from the beginning of the upstream operator chain. There is no mechanism to resume execution from a specific operator in the middle of a cold chain. The entire upstream is cold and re-executes from its cold source.
The contrast with Mono.just(value) is instructive: Mono.just(UUID.randomUUID()) evaluates the UUID at assembly time. The UUID is captured in the Mono.just’s closure as a concrete value. Re-subscribing to Mono.just(capturedValue) delivers the same capturedValue without re-executing any factory. But Mono.just(UUID.randomUUID()) at the top of a method that returns a Mono evaluates the UUID every time that method is called to assemble the chain, not every time the chain is subscribed. If the chain is assembled once per incoming request (which is the normal WebFlux handler pattern), and the chain is subscribed once per retry attempt (which is how retryWhen() works), then Mono.just(UUID.randomUUID()) gives a stable UUID across retries but a new UUID across requests — which is exactly the desired behavior.
The fix: generate UUID at assembly time using Mono.just(UUID.randomUUID()) or, better, compute a content-hash key from stable inputs outside the defer block and pass it into the chain as a captured value.
// Fix A: UUID at chain assembly time (Mono.just captures at assembly)
public Mono chargeCustomer(String customerId, String stripeCustomerId,
long amountCents, String billingPeriod) {
// UUID evaluated here — once per method call (once per incoming request).
// retryWhen() re-subscribes the chain but does NOT re-call this method.
// The Mono.just(key) below delivers the same key to every subscriber.
String key = UUID.randomUUID().toString();
return Mono.just(key)
.flatMap(k -> repository.save(
BillingAttempt.pending(customerId, amountCents, billingPeriod, k))
.thenReturn(k))
.flatMap(k -> stripeClient.createPaymentIntent(stripeCustomerId, amountCents, k))
.retryWhen(Retry.backoff(3, Duration.ofMillis(200))
.filter(ex -> ex instanceof WebClientResponseException.InternalServerError));
}
// Fix B: content-hash key — stable across requests for the same billing intent
public Mono chargeCustomer(String customerId, String stripeCustomerId,
long amountCents, String billingPeriod) {
// Content-hash: deterministic from domain inputs.
// If this method is called twice for the same customer + period + amount,
// the same key is produced — Stripe deduplicates the second request.
String key = UUID.nameUUIDFromBytes(
(customerId + ":" + billingPeriod + ":" + amountCents).getBytes())
.toString();
return Mono.just(key)
.flatMap(k -> repository.save(
BillingAttempt.pending(customerId, amountCents, billingPeriod, k))
.thenReturn(k))
.flatMap(k -> stripeClient.createPaymentIntent(stripeCustomerId, amountCents, k))
.retryWhen(Retry.backoff(3, Duration.ofMillis(200))
.filter(ex -> ex instanceof WebClientResponseException.InternalServerError));
}
Reactor Mono.fromCallable() variant: the same problem with a different operator name
Mono.fromCallable(() -> UUID.randomUUID().toString()) has identical semantics to Mono.defer(() -> Mono.just(UUID.randomUUID().toString())) with respect to re-subscription: the callable re-executes on every subscription, including each retryWhen() re-subscription. The failure mode is identical to Mode 1. fromCallable and defer both defer UUID evaluation to subscription time, and retryWhen() causes re-subscription. The fix is the same: move UUID generation outside the fromCallable block and capture it as a plain variable at assembly time.
// Also unsafe — fromCallable re-evaluates UUID on retryWhen() re-subscription:
return Mono.fromCallable(() -> UUID.randomUUID().toString())
.flatMap(key -> repository.save(...).thenReturn(key))
.flatMap(key -> stripeClient.createPaymentIntent(..., key))
.retryWhen(Retry.backoff(3, Duration.ofMillis(200)));
// Same fix: UUID outside Mono.fromCallable, captured as a local variable.
Mode 2 (Kotlin Flow): cold flow{} builder re-collection in a batch billing scenario
The developer is building a batch billing job in Kotlin using the Flow API. The job charges a list of customers reactively, with concurrency controlled via flatMapMerge. The Flow approach is chosen for its structured concurrency, backpressure, and the ability to apply retryWhen{} on individual failure modes. The billing flow emits pairs of (idempotency key, customer) from a flow{} builder, and the downstream flatMapMerge calls Stripe per customer.
The developer places UUID.randomUUID() inside the flow{} builder at the emit() site, reasoning that each customer should receive a fresh UUID at emission time — not at the time the flow graph is assembled. The retryWhen{} operator is placed at the end of the pipeline to retry on Stripe 429 or transient network errors.
The problem: the outer retryWhen{} re-collects the entire upstream cold flow on failure. Any failure in the downstream flatMapMerge — even for a single customer — causes retryWhen{} to signal a re-collection of the entire flow{} builder. The flow{} builder block re-executes from its beginning, emitting all customers again. UUID.randomUUID() generates UUID_B for every customer, including those who already received ch_A from the first collection pass. The blast radius is the entire batch: every customer charged in attempt 1 receives a second charge in attempt 2.
// BillingBatchService.kt — unsafe Mode 2: flow{} builder with outer retryWhen{}
@Service
class BillingBatchService(
private val repository: BillingAttemptRepository,
private val stripeClient: StripeWebClient
) {
// Developer's reasoning:
// "flow{} emits each customer lazily at collection time. UUID.randomUUID()
// inside the flow builder is natural — each customer gets a fresh UUID
// when the emission happens. retryWhen{} only retries the element that
// failed, not all upstream emissions."
//
// The problem:
// flow{} is a COLD builder. Every call to collect() re-executes the builder block.
// retryWhen{} on the outer flow causes re-collection from the BEGINNING of
// the upstream flow on any error. This means:
// - All customers are re-emitted.
// - UUID.randomUUID() generates UUID_B for all customers.
// - flatMapMerge calls Stripe for all customers again.
// - Customers who already received ch_A in attempt 1 receive ch_B in attempt 2.
// Blast radius = entire batch.
//
// Timeline:
// Collection 1 (attempt 1):
// flow{} builder emits (UUID_A1, cust1), (UUID_A2, cust2), ..., (UUID_AN, custN).
// flatMapMerge concurrently calls Stripe for all customers.
// Customers 1..K receive ch_A1..ch_AK (committed in Stripe).
// Customer K+1 Stripe call fails with 429 or 500.
// retryWhen{} catches the error — re-collects the upstream flow.
//
// Re-collection (attempt 2):
// flow{} builder re-executes from its first statement.
// emit(UUID.randomUUID() to cust1) → UUID_B1. ← NEW UUID for cust1
// emit(UUID.randomUUID() to cust2) → UUID_B2. ← NEW UUID for cust2
// ...
// flatMapMerge calls Stripe for ALL customers with UUID_B1..UUID_BN.
// Customers 1..K: ch_B1..ch_BK committed. ← DUPLICATE charges
fun chargeBatch(customers: List): Flow {
return flow {
for (customer in customers) {
// UUID generated per customer per emission — inside the cold builder.
val idempotencyKey = UUID.randomUUID().toString()
emit(idempotencyKey to customer)
}
}
.flatMapMerge(concurrency = 8) { (key, customer) ->
flow {
val result = stripeClient.createPaymentIntent(
customer.stripeCustomerId,
customer.amountCents,
key
)
repository.save(BillingAttempt.charged(customer.id, key, result.id))
emit(ChargeResult(customer.id, result.id, key))
}
}
.retryWhen { cause, attempt ->
// Developer expects this to retry only the failed customer's flow.
// Actually: retries re-collect the ENTIRE upstream cold flow.
attempt < 3 && (cause is StripeException || cause is IOException)
}
}
}
Why retryWhen{} in Kotlin Flow re-collects the entire upstream
Kotlin Flow’s retryWhen{} (and the simpler retry{}) are intermediate operators that act on the upstream flow. When the upstream flow emits an error terminal signal, retryWhen{} evaluates the predicate and, if it decides to retry, re-collects the upstream flow from its beginning. There is no mechanism in the Flow API to resume an upstream cold builder from a specific emission point. The cold builder is a function that starts from its first statement on each collect() call.
The developer’s misconception — that retryWhen{} retries only the failed element — comes from conflating Flow retry semantics with the semantics of per-element retry inside flatMapMerge. Per-element retry would require placing retryWhen{} inside the flatMapMerge lambda, wrapping only the inner flow for each customer. In that case, the inner flow for one customer re-collects only that customer’s inner flow — the outer flow{} builder that emits customers is not touched.
The distinction maps to Reactor: Flux.flatMap { innerFlux.retryWhen() } retries per element (each element’s inner Mono/Flux retries independently). Flux.flatMap { innerFlux }.retryWhen() retries the entire outer flux, which re-runs the cold source. The outer position of retryWhen determines its re-execution scope.
The fix: UUID outside the cold builder; per-element retry inside flatMapMerge
Two independent fixes are needed. First, UUID must be computed before the cold builder executes. If the list of customers is available before the flow starts, content-hash keys can be computed per customer outside the flow and passed as stable values into the emission. Second, per-element retry must be placed inside the flatMapMerge lambda so that only the failing customer’s inner flow retries, and the outer flow is not re-collected.
// Fix: content-hash keys pre-computed before flow builder; per-element retryWhen inside flatMapMerge
@Service
class BillingBatchService(
private val repository: BillingAttemptRepository,
private val stripeClient: StripeWebClient
) {
fun chargeBatch(customers: List): Flow {
// Step 1: Pre-compute content-hash keys for all customers BEFORE the flow starts.
// Keys are stable: if this method is called twice for the same batch,
// the same keys are produced and Stripe deduplicates.
val customerKeys: List> = customers.map { customer ->
val key = UUID.nameUUIDFromBytes(
"${customer.id}:${customer.billingPeriod}:${customer.amountCents}".toByteArray()
).toString()
key to customer
}
return flow {
// Emit pre-computed (key, customer) pairs — no UUID generation here.
for ((key, customer) in customerKeys) {
emit(key to customer)
}
}
.flatMapMerge(concurrency = 8) { (key, customer) ->
// Step 2: per-element retry INSIDE flatMapMerge.
// Only this customer's inner flow retries on failure.
// The outer flow{} builder is never re-collected.
flow {
val result = stripeClient.createPaymentIntent(
customer.stripeCustomerId,
customer.amountCents,
key // same key on every retry attempt for this customer
)
repository.save(BillingAttempt.charged(customer.id, key, result.id))
emit(ChargeResult(customer.id, result.id, key))
}.retryWhen { cause, attempt ->
attempt < 3 && (cause is StripeException || cause is IOException)
}
}
}
}
Kotlin channelFlow{}: the same problem applies
channelFlow{} creates a cold channel-based flow that allows concurrent send() operations from coroutines launched inside the builder. Its coldness with respect to re-collection is identical to flow{}: retryWhen{} on the outer flow re-collects the channelFlow{} block, re-running all coroutines launched inside it and regenerating any UUIDs computed within them. The fix is the same: pre-compute keys outside the builder and pass them as captured values. The per-element retry pattern also applies: place retryWhen{} on the inner flow produced by the flatMapMerge or flatMapLatest lambda, not on the outer channelFlow{}.
Mode 3 (both frameworks): TransactionalOperator (Reactor) vs. @Transactional suspend fun (Kotlin) — @Retryable outer AOP proxy causes UUID re-generation in both, through different mechanisms
This mode illustrates how the same root cause — @Retryable as the outer AOP proxy — manifests differently in Project Reactor (via Mono re-subscription triggered by proceed()) and Kotlin coroutines (via a new coroutine invocation created by proceed()). Understanding both mechanically helps developers using either framework avoid making the same mistake.
Mode 3a: Reactor — @Retryable on a method returning Mono with TransactionalOperator.transactional()
The developer uses TransactionalOperator to manage R2DBC transactions programmatically inside a reactive service method. The service method is annotated with @Retryable to handle transient failures. The method returns a Mono<String> that, when subscribed, starts an R2DBC transaction, saves a billing record, calls Stripe, and commits. The UUID is generated inside the transactional Mono using Mono.defer{} for reactive correctness.
The problem: @Retryable’s AOP proxy wraps the Java method itself, not the Mono it returns. When the method is called, it returns a Mono (assembles the chain). The caller subscribes to the Mono. The subscription executes the chain. If the chain fails (the Mono emits an error), the error propagates out of the subscription to whatever subscribed to the returned Mono. The @Retryable proxy does not intercept the error signal from the Mono — it only intercepts exceptions thrown synchronously by the method invocation.
This is a subtler variant of the problem. To make @Retryable actually retry the Stripe call, the developer must block on the Mono inside the method body, or the subscription and retry must happen in a context where exceptions propagate synchronously. In a Spring WebFlux or R2DBC context, the typical pattern is to subscribe on the method’s reactive pipeline itself and not use @Retryable at all — using retryWhen() as a Reactor operator instead.
However, some developers use @Retryable with reactive methods in Spring MVC-style services (where the reactive chain is subscribed to synchronously via .block()), or in a context where the framework subscribes to the returned Mono and unwraps exceptions. In those contexts, @Retryable’s proceed() re-invokes the method, which re-assembles the chain and returns a new Mono. The new Mono is then subscribed to again. If UUID.randomUUID() was inside a Mono.defer{} in the chain, the re-assembly either captures a new UUID at Mono.just(UUID.randomUUID()) call time (if the UUID is in the method body outside the Mono) or re-generates via defer on the new subscription.
The common failing pattern in a blocking-compatible context:
// ReactiveBillingService.java — unsafe Mode 3a: @Retryable on Mono-returning method
// This is the pattern that arises when developers use Spring MVC + R2DBC together,
// or when they call .block() to bridge from reactive to imperative.
@Service
public class ReactiveBillingService {
private final TransactionalOperator transactionalOperator;
private final R2dbcBillingRepository repository;
private final StripeWebClient stripeClient;
// Developer's reasoning:
// "@Retryable on the method retries the entire billing operation including the
// R2DBC transaction. TransactionalOperator.transactional() wraps the reactive
// chain as a unit — @Retryable retries this unit on failure. The UUID inside
// Mono.defer{} is reactive-correct — it defers UUID generation to subscription
// time within the transactional boundary."
//
// The problem when .block() is used or the method is called in a blocking context:
// @Retryable outer proxy calls the method → chain assembled → .block() subscribed.
// If Stripe throws, .block() propagates the exception synchronously.
// @Retryable catches the exception and calls proceed().
// proceed() calls the method again → NEW chain assembled.
// NEW Mono.defer{} factory in the new chain → UUID_B on the new subscription.
// ch_B committed alongside ch_A.
//
// Even without .block(): if the Mono is subscribed by the caller in a way that
// propagates errors synchronously (e.g., .toFuture().get(), or a blocking
// test via StepVerifier that re-subscribes), the same re-assembly issue arises.
@Retryable(
retryFor = {StripeException.class, WebClientRequestException.class},
maxAttempts = 3,
backoff = @Backoff(delay = 300, multiplier = 2.0)
)
public Mono chargeCustomer(String customerId, String stripeCustomerId,
long amountCents, String billingPeriod) {
return transactionalOperator.transactional(
Mono.defer(() -> {
// UUID inside defer — re-evaluated on each subscription.
String key = UUID.randomUUID().toString();
return Mono.just(key);
})
.flatMap(key -> repository.save(
BillingAttempt.pending(customerId, amountCents, billingPeriod, key))
.thenReturn(key))
.flatMap(key -> stripeClient.createPaymentIntent(
stripeCustomerId, amountCents, key))
);
}
}
Mode 3b: Kotlin — @Retryable on a @Transactional suspend fun
The Kotlin variant has a mechanically different execution path but the same UUID re-generation outcome. The developer annotates a suspend fun with both @Transactional and @Retryable. Spring’s @Transactional supports Kotlin coroutines through CoroutineTransactionManager: when a suspend fun annotated with @Transactional is invoked, the transaction manager obtains a connection and binds it to the coroutine context via TransactionContext. Repositories called within the function share this connection.
Spring Retry’s @Retryable does not have native coroutine support. It intercepts at the AOP proxy level, which in Kotlin means it intercepts the Java method that wraps the suspend fun. In practice, this means @Retryable’s proxy calls proceed() to re-invoke the method, which starts a new coroutine execution. Each proceed() invocation is a fresh coroutine: a new activation frame, new local variables, new UUID.randomUUID() call, UUID_B.
The developer misconception for Kotlin: “Spring’s @Transactional has coroutine support — it integrates with the coroutine context. I assume @Retryable is similarly coroutine-aware and retries by resuming the suspended coroutine rather than starting a new invocation.” This assumption is incorrect. @Retryable’s AnnotationAwareRetryOperationsInterceptor is a standard AOP interceptor. It calls proceed() on the MethodInvocation, which re-invokes the method. For a suspend fun, this creates a new coroutine (via the Kotlin coroutine machinery that wraps the suspend function as a Continuation). The new coroutine starts from line 1 of the function body. There is no state sharing between the original coroutine and the retry’s coroutine.
// BillingService.kt — unsafe Mode 3b: @Retryable + @Transactional on suspend fun
@Service
class BillingService(
private val repository: BillingAttemptRepository,
private val stripeClient: StripeWebClient
) {
// Developer's reasoning:
// "@Transactional is coroutine-aware in Spring — the transaction is bound to
// the coroutine context. @Retryable retries the billing operation. I assume
// Spring Retry is also coroutine-aware and retries by resuming the coroutine,
// so UUID generated at the start of the suspend fun is stable across retries."
//
// The problem:
// @Retryable (order MAX-5) is outer proxy to @Transactional (order MAX).
// Spring Retry is NOT coroutine-aware.
// proceed() creates a new coroutine invocation on each retry attempt.
// UUID.randomUUID() at the top of the function generates UUID_B on attempt 2.
// @Transactional opens a new coroutine transaction context for each new invocation.
//
// Attempt 1:
// @Retryable proxy intercepts → calls proceed() → @Transactional proxy intercepts
// → starts coroutine with new transaction T1 → function body starts from line 1
// → UUID_A = UUID.randomUUID()
// → repository.save(attempt with UUID_A) in T1
// → stripeClient.createPaymentIntent(..., UUID_A) → ch_A committed in Stripe
// → Stripe returns 429 → StripeException thrown
// → @Transactional rolls back T1 → exception propagates
// → @Retryable catches exception → calls proceed() again
//
// Attempt 2:
// proceed() → NEW coroutine invocation → NEW activation frame
// → UUID_B = UUID.randomUUID() ← NEW UUID
// → @Transactional opens NEW transaction T2
// → repository.save(attempt with UUID_B) in T2
// → stripeClient.createPaymentIntent(..., UUID_B) → ch_B committed ← DUPLICATE
@Transactional
@Retryable(
retryFor = [StripeException::class],
maxAttempts = 3,
backoff = Backoff(delay = 300, multiplier = 2.0)
)
suspend fun chargeCustomer(
customerId: String,
stripeCustomerId: String,
amountCents: Long,
billingPeriod: String
): String {
// UUID at function entry — re-generated on each @Retryable proceed() invocation.
val idempotencyKey = UUID.randomUUID().toString()
repository.save(
BillingAttempt.pending(customerId, amountCents, billingPeriod, idempotencyKey)
)
val result = stripeClient.createPaymentIntent(
stripeCustomerId, amountCents, idempotencyKey
)
repository.updateStatus(idempotencyKey, "CHARGED", result.id)
return result.id
}
}
The fix for Mode 3a (Reactor): UUID outside the Mono chain; retryWhen() instead of @Retryable
For Reactor-based services, the correct retry mechanism for a reactive chain is the Reactor operator retryWhen(), not the Spring AOP @Retryable. retryWhen() operates at the reactive layer — it is part of the Mono/Flux operator pipeline. When @Retryable is used on a method that returns a Mono, it only intercepts synchronous exceptions thrown by the method itself, not asynchronous errors emitted by the Mono. Using @Retryable with reactive methods is architecturally mismatched unless the method calls .block() or wraps the reactive pipeline in a way that turns reactive errors into synchronous exceptions.
The UUID fix in the Reactor context: compute UUID outside the Mono chain. If the method assembles a new chain per request (which is the normal pattern), the UUID should be a plain Java variable in the method body, captured by the lambda closures in the chain. The Mono.just(capturedKey) then delivers the same key on every subscription of the chain, including re-subscriptions caused by an inner retryWhen() operator within the chain.
// Fix for Mode 3a: UUID at method body scope; retryWhen() as Reactor operator (not @Retryable)
@Service
public class ReactiveBillingService {
private final TransactionalOperator transactionalOperator;
private final R2dbcBillingRepository repository;
private final StripeWebClient stripeClient;
// No @Retryable annotation — retry handled by Reactor's retryWhen() inside the chain.
public Mono chargeCustomer(String customerId, String stripeCustomerId,
long amountCents, String billingPeriod) {
// UUID computed at method body scope — same UUID for all subscriptions of this chain.
// A new UUID is computed per incoming request (per method call), but not per retry.
String key = UUID.nameUUIDFromBytes(
(customerId + ":" + billingPeriod + ":" + amountCents).getBytes()).toString();
return transactionalOperator.transactional(
Mono.just(key)
.flatMap(k -> repository.save(
BillingAttempt.pending(customerId, amountCents, billingPeriod, k))
.thenReturn(k))
.flatMap(k -> stripeClient.createPaymentIntent(stripeCustomerId, amountCents, k))
)
// retryWhen() as a Reactor operator on the transactional Mono.
// This is the correct way to retry reactive chains — not @Retryable on the method.
// Note: retryWhen() here re-subscribes the transactional Mono — it re-opens the
// R2DBC transaction on each retry. This is intentional: if the transaction rolled back,
// the DB write needs to be retried too. key is captured above and stable across
// all retry subscriptions.
.retryWhen(Retry.backoff(3, Duration.ofMillis(300))
.filter(ex -> ex instanceof StripeException
|| ex instanceof WebClientRequestException));
}
}
The fix for Mode 3b (Kotlin): UUID in a non-retried facade; @Retryable on the inner service method that receives key as a parameter
The fix pattern for Kotlin coroutines is the same as for Spring Data JPA Mode 1 (covered in the Spring Data JPA post): separate UUID computation into a non-@Retryable facade, and apply @Retryable only to an inner method that receives the key as a stable parameter. The inner method’s parameter value does not change between proceed() invocations — Spring Retry re-uses the same MethodInvocation arguments on each retry, so the key parameter received by the inner method is the same UUID on every attempt.
// Fix for Mode 3b: UUID in non-@Retryable facade; inner @Retryable method receives key as parameter
@Service
class BillingFacade(private val billingService: BillingService) {
// No @Retryable — computes the key and delegates.
suspend fun chargeCustomer(
customerId: String,
stripeCustomerId: String,
amountCents: Long,
billingPeriod: String
): String {
// Content-hash key: stable across any retry or duplicate call for same billing intent.
val key = UUID.nameUUIDFromBytes(
"$customerId:$billingPeriod:$amountCents".toByteArray()
).toString()
// Passes the stable key to the @Retryable inner method.
return billingService.chargeWithKey(customerId, stripeCustomerId, amountCents, billingPeriod, key)
}
}
@Service
class BillingService(
private val repository: BillingAttemptRepository,
private val stripeClient: StripeWebClient
) {
// @Retryable here: if proceed() is called on retry, the 'key' parameter
// received from the MethodInvocation is the SAME string value on every attempt.
// No UUID.randomUUID() inside this method body.
@Transactional
@Retryable(
retryFor = [StripeException::class],
maxAttempts = 3,
backoff = Backoff(delay = 300, multiplier = 2.0)
)
suspend fun chargeWithKey(
customerId: String,
stripeCustomerId: String,
amountCents: Long,
billingPeriod: String,
key: String // ← stable parameter, same value on every retry attempt
): String {
// No UUID generation here — key is received as a parameter.
repository.save(
BillingAttempt.pending(customerId, amountCents, billingPeriod, key)
)
val result = stripeClient.createPaymentIntent(
stripeCustomerId, amountCents, key
)
repository.updateStatus(key, "CHARGED", result.id)
return result.id
}
}
Cross-framework comparison: Reactor Mono.defer{} vs. Kotlin flow{} vs. AOP proxy layer
| Framework | Re-execution mechanism | Operator/proxy that causes re-execution | UUID position in the unsafe pattern | Developer misconception |
|---|---|---|---|---|
| Reactor (Mode 1) | Re-subscription of cold Mono — defer{} factory re-executes |
retryWhen() re-subscribes upstream |
Inside Mono.defer{} or Mono.fromCallable{} |
Mono.defer{} lazily wraps UUID — retryWhen() retries from the failed operator, not the deferred source |
| Kotlin Flow (Mode 2) | Re-collection of cold flow{} builder — builder block re-executes |
Outer retryWhen{} on the flow re-collects from the upstream builder |
Inside flow{} builder at emit() site |
retryWhen{} retries the failed element, not all upstream emissions in the cold builder |
| Reactor (Mode 3a) | Re-assembly and re-subscription of Mono — new chain built by proceed() |
@Retryable outer AOP proxy calls proceed() → method re-invoked → new Mono assembled and subscribed |
Inside Mono.defer{} in the chain, or Mono.just(UUID.randomUUID()) in method body outside the Mono |
TransactionalOperator.transactional() wraps the chain as a committed unit — @Retryable retries from after the transaction |
| Kotlin coroutines (Mode 3b) | New coroutine invocation — new activation frame re-executes suspend fun body |
@Retryable outer AOP proxy calls proceed() → new coroutine starts from line 1 |
Inside suspend fun body, before first repository.save() |
Spring’s coroutine @Transactional support implies @Retryable is coroutine-aware — retry resumes the existing coroutine |
| Mode | DB state after failure | Stripe state | Detection signal | Fix |
|---|---|---|---|---|
1: Reactor defer{} + retryWhen() |
Attempt 1 R2DBC transaction may have rolled back (no DB record for UUID_A); attempt 2 committed UUID_B | ch_A from attempt 1 (committed before Stripe failure); ch_B from attempt 2 — two charges | Stripe customer has two PaymentIntents for same billing period; one has no billing_attempt counterpart; Stripe idempotency conflict log absent (different keys) | UUID at assembly time outside defer{}; captured as Java/Kotlin variable in method scope; Mono.just(capturedKey) delivers same value per re-subscription |
2: Kotlin Flow flow{} + outer retryWhen{} |
Attempt 1 repository saves may have committed for some customers before the failure; attempt 2 re-inserts UUID_B rows — duplicate key violation on idempotency_key unique index if present; otherwise two rows per customer | ch_A per customer in attempt 1 batch (all customers processed before failure); ch_B per customer in attempt 2 batch — N duplicate charges | Stripe dashboard: multiple PaymentIntents per customer for the same billing period; application logs: duplicate key insert errors on billing_attempt table; batch audit: charge count = 2 × customer count | Pre-compute content-hash keys before flow builder; pass as captured values into emission; move retryWhen{} inside flatMapMerge lambda for per-element retry scope |
3a: Reactor TransactionalOperator + @Retryable |
Same as Mode 1 — per-attempt transaction; UUID_B row committed; UUID_A row absent or rolled back | ch_A and ch_B committed; no cross-key deduplication in Stripe | Stripe customer has two PaymentIntents; no billing_attempt row for UUID_A; application logs: @Retryable retry attempt log entry at same timestamp as duplicate Stripe charge |
UUID at method body scope, captured in lambda; use retryWhen() as Reactor operator rather than @Retryable on reactive-returning method |
3b: Kotlin @Transactional suspend fun + @Retryable |
Attempt 1 transaction rolled back (UUID_A row absent); attempt 2 committed UUID_B row | ch_A from attempt 1 (committed before StripeException); ch_B from attempt 2 — two charges |
Stripe has two PaymentIntents; billing_attempt has only UUID_B row; @Retryable retry log appears after initial failure |
Separate UUID computation into non-@Retryable facade; pass as parameter to @Retryable + @Transactional inner method |
Reactor Flux.generate() and Kotlin Flow.asFlow() on Sequence: the safe patterns
Not all reactive sources have the re-execution problem. Understanding which patterns are safe helps developers make correct choices when building billing pipelines.
Reactor Flux.fromIterable() with pre-computed list: if the list containing (key, customer) pairs is built before the Flux assembly, Flux.fromIterable(preComputedList) delivers the pre-computed elements on each re-subscription without re-executing any factory function. The key values are fixed. This is the Reactor equivalent of the pre-computed customerKeys list in the Kotlin fix above.
Kotlin What makes a pattern safe: the key must not be computed inside a re-executed factory or builder. Whether the UUID is inside a Testing reactive retry idempotency requires verifying that the same Three reactive Stripe idempotency failure modes across Project Reactor and Kotlin Flow, each driven by the same root cause — UUID inside a re-executed code boundary — but through different reactive mechanisms: The structural rule applies across both frameworks: any code that computes the Stripe idempotency key must execute exactly once, before all retried boundaries. In reactive code, a “retried boundary” includes any re-subscription trigger ( Related posts in this series: Spring Data R2DBC + Kotlin Coroutines + Stripe — Your agent uses a vault key. Keybrake enforces a daily USD cap, endpoint allowlist, and merchant scope before each call exits — then logs every request with parsed cost. One-click revoke. No key rotation required.List.asFlow(): a List converted to a Flow via asFlow() is a cold flow that re-emits the same list elements on each re-collection — but if the list contains pre-computed keys, those keys are the same on every re-collection. retryWhen{} on a list.asFlow().flatMapMerge { ... }` chain with pre-computed keys per element is safe: re-collection delivers the same (key, customer) pairs every time.
Mono.defer{} factory, a Mono.fromCallable{} callable, a flow{} builder block, a channelFlow{} block, or a method body that is re-invoked by @Retryable, the outcome is the same: a new UUID on each execution. The pattern is safe when the UUID is a concrete value captured in a closure or passed as a parameter before any re-executed boundary.JUnit 5 + WireMock test patterns
Idempotency-Key header is sent to Stripe on all retry attempts, regardless of how many times the reactive operator or AOP proxy re-executes the code path. WireMock’s scenario API provides a clean way to fail once and succeed on the second attempt, then assert on the captured request headers.Mode 1 test: Reactor
retryWhen() — same key on re-subscription// ReactiveBillingServiceTest.java — Mode 1
@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.NONE)
@AutoConfigureWireMock(port = 0)
class ReactiveBillingServiceTest {
@Autowired
ReactiveBillingService billingService;
@Autowired
WireMockServer wireMockServer;
@Autowired
R2dbcBillingRepository repository;
@Test
void chargeCustomer_retryWhen_sameKeyOnBothAttempts() {
String successBody = """
{"id":"pi_ok","object":"payment_intent","status":"succeeded","amount":5000}
""";
// First call: Stripe returns 500 (triggers retryWhen).
wireMockServer.stubFor(post(urlEqualTo("/v1/payment_intents"))
.inScenario("stripe-retry")
.whenScenarioStateIs(Scenario.STARTED)
.willReturn(aResponse().withStatus(500)
.withBody("{\"error\":{\"type\":\"api_error\"}}"))
.willSetStateTo("attempt-2"));
// Second call: Stripe succeeds.
wireMockServer.stubFor(post(urlEqualTo("/v1/payment_intents"))
.inScenario("stripe-retry")
.whenScenarioStateIs("attempt-2")
.willReturn(aResponse().withStatus(200)
.withHeader("Content-Type", "application/json")
.withBody(successBody)));
// Act: subscribe and block (or use StepVerifier).
StepVerifier.create(billingService.chargeCustomer(
"cust_1", "stripe_cus_1", 5000L, "2026-Q4"))
.expectNextCount(1)
.verifyComplete();
// Assert: Stripe received exactly 2 requests, both with the SAME Idempotency-Key.
ListMode 2 test: Kotlin Flow batch — per-element retry inside
flatMapMerge, same key per customer// BillingBatchServiceTest.kt — Mode 2
@SpringBootTest
@AutoConfigureWireMock(port = 0)
class BillingBatchServiceTest {
@Autowired
lateinit var billingBatchService: BillingBatchService
@Autowired
lateinit var wireMockServer: WireMockServer
@Autowired
lateinit var billingAttemptRepository: BillingAttemptRepository
@Test
fun chargeBatch_perElementRetry_sameKeyPerCustomer() = runTest {
val customers = listOf(
Customer("cust-1", "stripe_cus_1", 1000L, "2026-Q4"),
Customer("cust-2", "stripe_cus_2", 2000L, "2026-Q4")
)
// Pre-compute expected keys (same algorithm as the service).
val expectedKeys = customers.map { c ->
UUID.nameUUIDFromBytes("${c.id}:${c.billingPeriod}:${c.amountCents}".toByteArray()).toString()
}
// Stub: first call per customer fails, second succeeds.
// Use customer-specific stubs differentiated by request body.
for ((i, key) in expectedKeys.withIndex()) {
wireMockServer.stubFor(post(urlEqualTo("/v1/payment_intents"))
.withHeader("Idempotency-Key", equalTo(key))
.inScenario("retry-cust-$i")
.whenScenarioStateIs(Scenario.STARTED)
.willReturn(aResponse().withStatus(429)
.withBody("{\"error\":{\"type\":\"rate_limit_error\"}}"))
.willSetStateTo("retry"))
wireMockServer.stubFor(post(urlEqualTo("/v1/payment_intents"))
.withHeader("Idempotency-Key", equalTo(key))
.inScenario("retry-cust-$i")
.whenScenarioStateIs("retry")
.willReturn(aResponse().withStatus(200)
.withHeader("Content-Type", "application/json")
.withBody("""{"id":"pi_$i","status":"succeeded","amount":${customers[i].amountCents}}""")))
}
// Act: collect all results.
val results = billingBatchService.chargeBatch(customers).toList()
// Assert: 2 results, one per customer.
assertThat(results).hasSize(2)
// Assert: each customer's Stripe calls used the same key across retry attempts.
for (key in expectedKeys) {
val requests = wireMockServer.findAll(
postRequestedFor(urlEqualTo("/v1/payment_intents"))
.withHeader("Idempotency-Key", equalTo(key)))
// Exactly 2 calls per customer (1 failure + 1 success) with the same key.
assertThat(requests).hasSize(2)
}
// Assert: 2 billing_attempt rows, one per customer, no duplicates.
assertThat(billingAttemptRepository.count()).isEqualTo(2)
}
}
Mode 3b test: Kotlin
@Retryable facade — same key parameter on each proceed()// BillingFacadeTest.kt — Mode 3b
@SpringBootTest
@AutoConfigureWireMock(port = 0)
class BillingFacadeTest {
@Autowired
lateinit var billingFacade: BillingFacade
@Autowired
lateinit var wireMockServer: WireMockServer
@Autowired
lateinit var billingAttemptRepository: BillingAttemptRepository
@Test
fun chargeCustomer_retryableWithStableKey_stripeSeesOneKey() = runTest {
val customerId = "cust-facade-test"
val billingPeriod = "2026-Q4"
val amountCents = 3000L
// Expected content-hash key (same algorithm as the facade).
val expectedKey = UUID.nameUUIDFromBytes(
"$customerId:$billingPeriod:$amountCents".toByteArray()
).toString()
// First call: Stripe returns 429 (triggers @Retryable).
wireMockServer.stubFor(post(urlEqualTo("/v1/payment_intents"))
.withHeader("Idempotency-Key", equalTo(expectedKey))
.inScenario("facade-retry")
.whenScenarioStateIs(Scenario.STARTED)
.willReturn(aResponse().withStatus(429)
.withBody("{\"error\":{\"type\":\"rate_limit_error\"}}"))
.willSetStateTo("attempt-2"))
wireMockServer.stubFor(post(urlEqualTo("/v1/payment_intents"))
.withHeader("Idempotency-Key", equalTo(expectedKey))
.inScenario("facade-retry")
.whenScenarioStateIs("attempt-2")
.willReturn(aResponse().withStatus(200)
.withHeader("Content-Type", "application/json")
.withBody("""{"id":"pi_facade_ok","status":"succeeded","amount":$amountCents}""")))
// Act.
val piId = billingFacade.chargeCustomer(
customerId, "stripe_cus_facade", amountCents, billingPeriod)
assertThat(piId).isEqualTo("pi_facade_ok")
// Assert: Stripe saw exactly 2 calls, both with the same expected key.
val requests = wireMockServer.findAll(
postRequestedFor(urlEqualTo("/v1/payment_intents"))
.withHeader("Idempotency-Key", equalTo(expectedKey)))
assertThat(requests).hasSize(2)
// Assert: one billing_attempt row with the expected key.
val attempts = billingAttemptRepository.findByCustomerId(customerId)
assertThat(attempts).hasSize(1)
assertThat(attempts.first().idempotencyKey).isEqualTo(expectedKey)
}
}
Summary
Mono.defer{} inside retryWhen(): Mono.defer{} defers evaluation to subscription time. retryWhen() re-subscribes the entire upstream chain, causing defer{}’s factory to re-execute. UUID_B on the second subscription. The same applies to Mono.fromCallable{}. Fix: UUID at method body scope (assembly time), captured as a plain variable into the reactive chain. Use Mono.just(capturedKey) or content-hash key computed before chain assembly.flow{} cold builder + outer retryWhen{}: the flow{} builder block re-executes on every re-collection, including each retryWhen{} retry. UUID inside emit() regenerates for every customer in the batch. Blast radius = entire batch. Fix: pre-compute content-hash keys for all customers before the flow starts, pass as captured values in the emission. Move retryWhen{} inside flatMapMerge lambda for per-element retry scope.@Retryable outer AOP proxy with TransactionalOperator (Reactor) or @Transactional suspend fun (Kotlin): proceed() re-invokes the method — in Reactor, this assembles and subscribes a new chain; in Kotlin, this creates a new coroutine invocation. UUID at the beginning of the method/chain is re-generated in both. Developer misconception that TransactionalOperator is a committed unit or that Spring Retry is coroutine-aware. Fix for Reactor: use retryWhen() as a Reactor operator on the chain, not @Retryable on the reactive-returning method. Fix for Kotlin: separate UUID computation into a non-@Retryable facade; apply @Retryable to an inner method receiving the key as a stable parameter.retryWhen(), retry()), any re-collection trigger (Flow.retryWhen{}), and any AOP proceed() call from @Retryable. Content-hash keys computed from stable domain inputs (customer ID, billing period, amount) are more robust than random UUIDs because they remain stable even across duplicate method calls from upstream schedulers or HTTP clients.@Transactional suspend fun + @Retryable fresh coroutine, flatMap placement of retryWhen(), TransactionalOperator wrapping UUID generation in the same framework. Spring WebFlux Functional Endpoints + Stripe — bodyToMono().cache(), WebClient header customizer, Mono.fromCallable() in per-item batch retry. Kotlin Flow + Exposed + Stripe — Exposed’s built-in deadlock retry re-executing the transaction block, Flow.retry{} cold builder re-collection, flatMapLatest cancellation leaving committed Stripe charges.Keybrake: per-call enforcement for your agent’s Stripe calls