Spring Data R2DBC, Kotlin Coroutines, and Stripe Integration: How @Transactional suspend fun + @Retryable Starts a Fresh Coroutine on Each Retry, flatMap Placement of retryWhen() Restarts All Customer Pipelines and Regenerates Keys, and TransactionalOperator.transactional() Wrapping UUID.randomUUID() Creates UUID_B on Each Retry Subscription
Spring Data R2DBC with Kotlin coroutines is one of the most expressive reactive stacks on the JVM — and one of the most treacherous surfaces for Stripe idempotency keys. Three independent mechanisms all converge on the same outcome: UUID.randomUUID() runs twice, Stripe sees a new key, and a customer receives a second charge.
This post covers three failure modes specific to the R2DBC + Kotlin coroutine combination. They are structurally distinct from the Kotlin coroutines + @Transactional + @Retryable post (which uses JDBC with blocking coroutine bridges), the Kotlin Flow + Exposed post (which focuses on cold Flow re-collection semantics), and the Spring WebFlux post (which uses Java reactive operators without coroutine bridging). The modes here are specific to how Spring’s AOP proxy interacts with R2DBC reactive contexts, how Reactor’s retryWhen() operator re-subscribes upstream publishers, and how TransactionalOperator.transactional() wraps Mono compositions that include UUID generation.
Background: R2DBC transactions, reactive contexts, and the coroutine bridge
Spring Data R2DBC uses the R2dbcTransactionManager to manage transactions reactively. Unlike JDBC, which binds the transaction to a thread via ThreadLocal, R2DBC binds the transaction (and the underlying database connection) to the Reactor Context — an immutable key-value map that flows downstream through a reactive pipeline via Subscriber.currentContext(). The transaction is stored in the TransactionSynchronizationManager keyed to the reactive context, not to a thread.
When you write a Spring Data R2DBC service in Kotlin using coroutines, the coroutine bridge (CoroutinesUtils.invokeSuspendingFunction) converts a suspend fun into a Reactor Mono. The @Transactional annotation on a suspend fun is handled by CoroutinesTransactionInterceptor in Spring Framework 5.3+, which wraps the resulting Mono in a TransactionalOperator.executeAndAwait() call, binding the R2DBC transaction to the new reactive context of that subscription.
The key consequence: each new subscription to a reactive publisher triggers a new reactive context, a new transaction, and a new R2DBC connection from the pool. Whatever triggers a re-subscription — AOP proceed(), retryWhen(), or explicit re-subscription — effectively resets everything associated with the reactive context, including the R2DBC connection. Local variables computed inside the re-subscribed publisher, including UUID.randomUUID() calls, execute again with fresh bindings.
Stripe’s idempotency key contract: the Idempotency-Key header deduplicates requests by key within 24 hours for the same endpoint. Two requests with different keys for the same customer at the same amount are two charges, not one. The contract is per-key, not per-customer, not per-amount.
Mode 1: @Transactional suspend fun + @Retryable — AOP proceed() creates a fresh coroutine and a new reactive context — UUID_B on retry
The stacking of @Transactional and @Retryable on a suspend fun creates an AOP interceptor chain. Spring Framework’s CoroutinesTransactionInterceptor handles the @Transactional side: it invokes the next interceptor (or the target method), receives a Deferred result, and wraps the underlying Mono in a reactive transaction. Spring Retry’s AnnotationAwareRetryOperationsInterceptor handles the @Retryable side: it wraps the entire interceptor chain in a RetryTemplate.
The ordering of these two interceptors determines which one is outer. By default, Spring Retry’s advisor has @Order(Ordered.LOWEST_PRECEDENCE - 1) and Spring’s transaction advisor has @Order(Ordered.LOWEST_PRECEDENCE - 10). Lower precedence number = higher priority = outer. The transaction interceptor is inner; the retry interceptor is outer. This means the retry loop wraps the transactional call: on retry, RetryTemplate calls proceed() on the interceptor chain, which re-enters CoroutinesTransactionInterceptor, which opens a new transaction, which dispatches a new coroutine.
Every time proceed() is called, CoroutinesUtils.invokeSuspendingFunction creates a new Mono from the suspend fun body. A new subscription to that Mono starts. A new reactive context is created for the subscription. R2dbcTransactionManager obtains a new connection from the pool. The suspend fun body executes from its first statement with fresh local variable bindings. UUID.randomUUID() inside the body generates a new UUID — UUID_B on the second attempt.
// BillingService.kt — unsafe mode 1: UUID inside @Transactional @Retryable suspend fun
@Service
class BillingService(
private val stripeClient: StripeClient,
private val billingRepository: BillingAttemptRepository
) {
// Developer's reasoning:
// "@Transactional on a suspend fun uses CoroutinesTransactionInterceptor,
// which is coroutine-aware — it stores the transaction in the coroutine context.
// @Retryable is also a Spring annotation — maybe Spring's retry support for
// suspend functions is coroutine-aware too, and the retry re-enters the existing
// reactive context, meaning UUID.randomUUID() doesn't re-execute because the
// 'current transaction context' carries over the already-computed key."
//
// The problem: @Retryable is a Java AOP framework (not coroutine-aware).
// proceed() calls CoroutinesTransactionInterceptor fresh.
// A new Mono is created from the suspend fun body.
// A new reactive context is created for the new subscription.
// UUID.randomUUID() generates UUID_B inside the new subscription.
// Attempt 1: UUID_A → Stripe commits ch_A → StripeConnectException thrown.
// Attempt 2: proceed() → new Mono → new reactive context → new R2DBC connection.
// UUID_B generated → Stripe commits ch_B → two charges, one customer.
@Retryable(
retryFor = [StripeConnectException::class],
maxAttempts = 3,
backoff = Backoff(delay = 500, multiplier = 2.0)
)
@Transactional
suspend fun chargeCustomer(
customerId: String,
amountCents: Long,
billingPeriod: String
): String {
// UUID computed here — inside the suspend fun body.
// On every proceed() call: new Mono, new reactive context, UUID_B.
val idempotencyKey = UUID.randomUUID().toString()
val params = PaymentIntentCreateParams.builder()
.setAmount(amountCents)
.setCurrency("usd")
.setCustomer(customerId)
.setConfirm(true)
.build()
val pi = withContext(Dispatchers.IO) {
stripeClient.paymentIntents().create(
params,
RequestOptions.builder().setIdempotencyKey(idempotencyKey).build()
)
}
billingRepository.save(BillingAttempt(
customerId = customerId,
stripePaymentIntentId = pi.id,
billingPeriod = billingPeriod,
amountCents = amountCents,
idempotencyKey = idempotencyKey
))
return pi.id
}
}
Why the misconception arises from Spring’s coroutine integration
The developer’s mental model is grounded in a real feature of Spring Framework: CoroutinesTransactionInterceptor is genuinely coroutine-aware. It integrates with the Kotlin coroutine system via reactor.kotlin.core: it invokes the suspend function through CoroutinesUtils.invokeSuspendingFunction, which bridges the coroutine to a Reactor Mono, and then applies the reactive transaction via TransactionalOperator. This is qualitatively different from how Spring’s original TransactionInterceptor works for JDBC with ThreadLocal bindings.
The mistake is extrapolating from “Spring has coroutine-aware transaction support” to “Spring Retry’s @Retryable might also be coroutine-aware.” Spring Retry predates coroutine support by many years. The AnnotationAwareRetryOperationsInterceptor uses proceed() on a plain MethodInvocation — a raw Java reflection call. It is not aware of the Kotlin coroutine dispatch model, the reactive context, or the R2DBC transaction manager. It calls proceed() and the Java method invocation infrastructure hands a new coroutine to CoroutinesUtils.invokeSuspendingFunction. From that point on, everything is fresh: new coroutine, new Mono, new reactive context, new R2DBC connection, new local variable bindings.
There is an additional subtlety for R2DBC specifically: the reactive context that carries the transaction is created per-subscription. The first subscription to the Mono (attempt 1) creates context A, which carries transaction A and connection A. When that subscription terminates (the suspend fun throws and the Mono signals onError), context A is abandoned. The retry interceptor calls proceed() again. A new subscription (attempt 2) creates context B, carrying transaction B and connection B. There is no mechanism by which context B inherits anything from context A. The UUID generation is entirely local to the method body and is re-run fresh in context B.
The fix: content-hash key computed in the caller before the @Retryable boundary
Pass the idempotency key as a parameter to the @Transactional @Retryable suspend fun. The key is computed once in the non-retried caller and passed through all proceed() invocations as a stable parameter value.
// BillingService.kt — fixed mode 1: key passed as parameter
@Service
class BillingService(
private val stripeClient: StripeClient,
private val billingRepository: BillingAttemptRepository
) {
@Retryable(
retryFor = [StripeConnectException::class],
maxAttempts = 3,
backoff = Backoff(delay = 500, multiplier = 2.0)
)
@Transactional
suspend fun chargeCustomer(
customerId: String,
amountCents: Long,
billingPeriod: String,
idempotencyKey: String // ← stable parameter from caller
): String {
// idempotencyKey is a parameter — not computed inside this function.
// proceed() on attempt 2 passes the same idempotencyKey value.
// Stripe deduplicates: attempt 2 returns the cached result of attempt 1.
// No ch_B.
val params = PaymentIntentCreateParams.builder()
.setAmount(amountCents)
.setCurrency("usd")
.setCustomer(customerId)
.setConfirm(true)
.build()
val pi = withContext(Dispatchers.IO) {
stripeClient.paymentIntents().create(
params,
RequestOptions.builder().setIdempotencyKey(idempotencyKey).build()
)
}
billingRepository.save(BillingAttempt(
customerId = customerId,
stripePaymentIntentId = pi.id,
billingPeriod = billingPeriod,
amountCents = amountCents,
idempotencyKey = idempotencyKey
))
return pi.id
}
}
// BillingOrchestrator.kt — key computed before the @Retryable boundary
@Service
class BillingOrchestrator(
private val billingService: BillingService
) {
suspend fun runBillingCycle(customers: List, billingPeriod: String) {
customers.forEach { customer ->
// Content-hash key computed here — outside the @Retryable proxy.
// Same customer + period + amount = same key across all proceed() retries.
val key = contentHashKey(customer.id, billingPeriod, customer.amountCents)
billingService.chargeCustomer(
customerId = customer.id,
amountCents = customer.amountCents,
billingPeriod = billingPeriod,
idempotencyKey = key
)
}
}
private fun contentHashKey(customerId: String, period: String, amountCents: Long): String {
val input = "$customerId|$period|$amountCents"
val digest = MessageDigest.getInstance("SHA-256").digest(
input.toByteArray(Charsets.UTF_8)
)
return digest.joinToString("") { "%02x".format(it) }
}
}
Why the alternative of putting @Retryable inside @Transactional does not help
Some developers attempt to reverse the AOP order by using @Order to make @Transactional outer and @Retryable inner. The idea: wrap a single transaction around all retry attempts so that a partial retry does not escape the transaction boundary. For JDBC with ThreadLocal binding, this avoids the “new transaction on retry” problem. For R2DBC, it creates a different problem: the outer transaction’s reactive context is created for the first subscription. If the inner retry calls proceed() a second time, the new Mono from the suspend fun body gets its own new reactive context. The outer TransactionalOperator cannot propagate its context into the inner Mono’s subscription context unless the propagation is done explicitly via Context.of() or mono { }.contextWrite(outerContext). Spring’s standard AOP interceptor chain does not do this. The result is that the inner Mono (attempt 2) does not participate in the outer transaction and UUID.randomUUID() still generates UUID_B.
The content-hash key passed as a parameter is the only reliable fix regardless of AOP order.
Mode 2: Flux.flatMap { val key = UUID.randomUUID()… }.retryWhen() outside flatMap — entire customer pipeline restarts, UUID_B per customer
The second failure mode is specific to batch billing pipelines written with Reactor’s Flux operators and Kotlin coroutine integration. The developer intends retryWhen() to retry only the failed element’s processing. The actual behavior of retryWhen() placed after flatMap is to re-subscribe the entire upstream Flux — re-executing the flatMap body for every element in the customer list, not just the one that failed.
// BillingPipeline.kt — unsafe mode 2: retryWhen() OUTSIDE flatMap
@Service
class BillingPipeline(
private val billingRepository: BillingAttemptRepository,
private val stripeService: ReactiveStripeService
) {
fun runBillingFlux(customerFlux: Flux<Customer>): Flux<String> {
return customerFlux
.flatMap { customer ->
// UUID computed inside flatMap — re-executes on retryWhen re-subscription.
// Developer's reasoning:
// "flatMap processes each customer independently.
// retryWhen() after flatMap will retry only the element
// that failed — similar to how a per-element error handler works.
// The failed customer's flatMap body re-executes with a new UUID.
// Successfully processed customers are downstream of flatMap
// and are not re-processed."
//
// The problem: retryWhen() re-subscribes the ENTIRE upstream Flux
// (including customerFlux and the flatMap operator).
// ALL customers re-enter the flatMap body — not just the failed one.
// UUID.randomUUID() generates UUID_B for every customer.
// billingRepository.save() inserts a new row for every customer (UUID_B).
// stripeService.charge() sends UUID_B for every customer —
// ch_B per customer whose UUID_A charge committed in the first subscription.
val key = UUID.randomUUID().toString() // ← re-runs per re-subscription
billingRepository.save(
BillingAttempt(customerId = customer.id, stripeKey = key)
).flatMap { _ ->
stripeService.charge(key, customer.amountCents, customer.id)
}
}
.retryWhen( // ← OUTSIDE flatMap: re-subscribes customerFlux
Retry.backoff(3, Duration.ofSeconds(1))
.filter { it is StripeConnectException }
)
}
}
When stripeService.charge() throws for customer C₇ and the StripeConnectException propagates out of the flatMap operator, retryWhen() receives an onError signal. It re-subscribes from the beginning of the upstream publisher — which is customerFlux.flatMap { ... }. The reactive subscription starts over. customerFlux begins emitting customers again from C₁. The flatMap body re-executes for C₁: new UUID.randomUUID() call — UUID_B for C₁ — new billingRepository.save() — new stripeService.charge(UUID_B) — ch_B for C₁, who had already received ch_A_C₁ in the first subscription before C₇ failed.
This behavior is not a bug in Reactor. It is the documented behavior of retryWhen(): re-subscribe the upstream. The problem is the developer’s mental model of what “upstream” means in this context.
The upstream boundary of retryWhen()
In Reactor, a publisher chain is a description of a data flow. Operators are composited: each operator wraps the previous publisher in a new publisher. When you write customerFlux.flatMap { ... }.retryWhen(...), the chain is:
FluxFlatMapwrapscustomerFlux(the source)FluxRetryWhenwrapsFluxFlatMap
FluxRetryWhen subscribes to FluxFlatMap. When FluxFlatMap signals onError, FluxRetryWhen re-subscribes to FluxFlatMap. FluxFlatMap’s subscription to customerFlux is cancelled and a new subscription is started. customerFlux begins emitting elements again from the start. The flatMap function lambda — which contains UUID.randomUUID() — is invoked for each newly emitted element.
There is no per-element tracking in FluxRetryWhen. It does not know which element caused the onError signal, nor does it try to re-process only that element. The re-subscription is a full restart of the wrapped publisher.
Why R2DBC save() inserts a new row on re-subscription
Spring Data R2DBC’s ReactiveCrudRepository.save() uses the entity’s @Id field to determine whether to issue an INSERT or an UPDATE. The default isNew() detection (via EntityInformation<T, ID>) returns true if the ID is null, triggering an INSERT; it returns false if the ID is non-null, triggering an UPDATE (actually a SELECT then either INSERT or UPDATE depending on the implementation).
In mode 2, a fresh BillingAttempt(...) is constructed inside the flatMap body on each re-subscription. The fresh entity has a null ID (if using @GeneratedValue). save() treats it as a new entity and issues an INSERT. The first subscription produced row R₁ with stripe_key = UUID_A. The second subscription (after retryWhen re-subscribed) produces row R₂ with stripe_key = UUID_B. Both rows exist in the database alongside ch_A and ch_B in Stripe.
The developer might add a unique constraint on (customer_id, billing_period) as a safety net. This constraint violation will be thrown during the second save() — before stripeService.charge(UUID_B) executes. However, in some failure modes the Stripe charge has already been sent before the DB write (if the pipeline order is charge → save rather than save → charge), or the unique constraint protects only the billing period level, not the attempt level. The content-hash key is the correct fix regardless.
The fix: move retryWhen() inside the flatMap body for per-element retry
For per-element retry semantics, retryWhen() must be applied to the per-element publisher — inside the flatMap lambda — not to the outer Flux. Additionally, the idempotency key must be computed outside the retried operator chain.
// BillingPipeline.kt — fixed mode 2: retryWhen() INSIDE flatMap + key outside retried chain
@Service
class BillingPipeline(
private val billingRepository: BillingAttemptRepository,
private val stripeService: ReactiveStripeService
) {
fun runBillingFlux(customerFlux: Flux<Customer>): Flux<String> {
return customerFlux
.flatMap { customer ->
// Key computed outside the retried chain.
// Same customer + billing period + amount = same key across retries.
// retryWhen() is inside the flatMap — it re-subscribes only this
// customer's Mono, not the entire customerFlux.
val key = contentHashKey(customer.id, customer.billingPeriod, customer.amountCents)
billingRepository.save(
BillingAttempt(customerId = customer.id, stripeKey = key)
).flatMap { _ ->
stripeService.charge(key, customer.amountCents, customer.id)
}.retryWhen( // ← INSIDE flatMap: re-subscribes only this customer's Mono
Retry.backoff(3, Duration.ofSeconds(1))
.filter { it is StripeConnectException }
)
}
}
private fun contentHashKey(customerId: String, period: String, amountCents: Long): String {
val input = "$customerId|$period|$amountCents"
val digest = MessageDigest.getInstance("SHA-256").digest(
input.toByteArray(Charsets.UTF_8)
)
return digest.joinToString("") { "%02x".format(it) }
}
}
With retryWhen() inside the flatMap, the re-subscription boundary is the per-customer Mono chain (save().flatMap { charge() }), not customerFlux. Only the failing customer’s Mono is re-subscribed. The key variable is captured from the outer flatMap lambda scope (computed once per customer element) and is not re-evaluated on retryWhen re-subscription because retryWhen does not cross the flatMap boundary.
Additionally, with a content-hash key computed from customer.billingPeriod, even if retryWhen() were accidentally placed outside the flatMap in a future refactor, the key would remain stable for each customer because it is deterministic from the input data. The content-hash key provides defense-in-depth against future accidental re-subscription.
A note on flatMapSequential and concurrency
Using flatMap for batch billing processes all customers concurrently (up to the default concurrency limit of 256). If Stripe rate limits or per-account concurrency limits are a concern, use flatMapSequential (sequential processing, preserves order) or flatMap(customer -> ..., concurrency)) with an explicit concurrency cap. Neither flatMapSequential nor a reduced concurrency cap changes the retryWhen placement rules: the same re-subscription semantics apply regardless of concurrency settings. The key must still be outside the retried operator chain, and retryWhen must still be inside the flatMap for per-element semantics.
Mode 3: TransactionalOperator.transactional(Mono.fromCallable { UUID.randomUUID() }…).retryWhen() — re-subscription to transactional() re-runs UUID generation inside the new transaction
The third failure mode arises when developers use TransactionalOperator.transactional() explicitly (rather than @Transactional on a method) to wrap a Mono composition, and place retryWhen() after the transactional() operator. The intent is: the transactional() block commits the DB write and the Stripe charge atomically (or as close to atomically as possible), and retryWhen() retries the whole unit on network failure. The problem: retryWhen() after transactional() re-subscribes the entire transactional() Mono. The UUID generation inside the Mono re-executes in the new subscription with a new reactive context, producing UUID_B.
// BillingService.kt — unsafe mode 3: UUID inside TransactionalOperator.transactional()
@Service
class BillingService(
private val transactionalOperator: TransactionalOperator,
private val billingRepository: BillingAttemptRepository,
private val stripeClient: ReactiveStripeService
) {
fun chargeCustomer(customer: Customer): Mono<String> {
// Developer's intent:
// "I want the DB write and the Stripe charge to happen atomically.
// I'll wrap both in TransactionalOperator.transactional().
// If there's a network failure, retryWhen() will retry the whole unit.
// The transactional() block already committed the DB write with UUID_A;
// when retryWhen() re-tries, it will see the committed record and
// — somehow — use the same UUID_A for the Stripe charge."
//
// The problem: retryWhen() re-subscribes the transactional() Mono.
// A new subscription to Mono.fromCallable { UUID.randomUUID() } starts.
// UUID.randomUUID() runs again inside the new subscription — UUID_B.
// The new transaction commits billingRepository.save() with UUID_B (new row).
// stripeClient.charge(UUID_B) sends ch_B.
// Two DB rows (UUID_A, UUID_B) and two Stripe charges (ch_A, ch_B).
return transactionalOperator.transactional(
Mono.fromCallable { UUID.randomUUID().toString() } // ← re-runs on re-subscription
.flatMap { key ->
billingRepository.save(
BillingAttempt(
customerId = customer.id,
stripeKey = key,
billingPeriod = customer.billingPeriod
)
).then(stripeClient.charge(key, customer.amountCents, customer.id))
}
).retryWhen( // ← re-subscribes transactional() Mono on failure
Retry.backoff(3, Duration.ofSeconds(1))
.filter { it is StripeConnectException }
)
}
}
What actually happens on retryWhen re-subscription
The Mono chain assembled by transactionalOperator.transactional(...) is a cold publisher. “Cold” means that each new subscriber gets a new execution: a new reactive context, a new transaction, new local state. TransactionalOperator.transactional(mono) does not cache or memoize the result of the first subscription. It wraps the passed mono in a new Mono that, on each new subscription, begins a new transaction using the R2dbcTransactionManager.
When retryWhen() re-subscribes after a failure, it subscribes to the TransactionalMono wrapper again. The TransactionalMono begins a new transaction. Inside the new transaction, the inner Mono.fromCallable { UUID.randomUUID() } is subscribed fresh. Mono.fromCallable executes the Callable lazily — on subscription — so UUID.randomUUID() runs on this second subscription, producing UUID_B. The flatMap body executes with UUID_B. billingRepository.save() creates a new BillingAttempt with UUID_B (new entity, null ID, treated as new by R2DBC, new INSERT). The transaction commits with the new row. stripeClient.charge(UUID_B) creates ch_B.
The developer’s misconception is that transactional() is a “committed unit” from which retryWhen() retries only the tail (the Stripe charge). This would be true if retryWhen() were placed inside the transactional() block, after the Stripe charge operator. In that position, a retry re-subscribes from the Stripe charge operator onward — but even then, only if the UUID is computed before the retryWhen() boundary. Placing retryWhen() outside the transactional() block means retrying the entire transactional unit, including the UUID generation at its beginning.
Understanding Mono.fromCallable vs. Mono.just for UUID generation
An important distinction: Mono.just(UUID.randomUUID().toString()) does not have this problem. Mono.just(value) evaluates value eagerly at construction time — the UUID is computed once when the just operator is created — and the same value is emitted on each subscription. Mono.fromCallable { UUID.randomUUID().toString() } evaluates the lambda lazily on each subscription, producing a new UUID each time.
Both forms appear in production code. Mono.fromCallable is commonly used for blocking operations (R2DBC queries, SHA-256 computation) that should not run at assembly time. If UUID generation is bundled into a fromCallable alongside a legitimately blocking operation (e.g., Mono.fromCallable { computeHashKey(); callBlockingApi() }), the UUID generation inherits the deferred execution of the blocking operation and re-runs on re-subscription.
// The distinction between eager and deferred UUID generation in Reactor:
// SAFE: eager — UUID computed once at assembly time
val key1 = UUID.randomUUID().toString() // computed NOW
val mono1 = Mono.just(key1) // same key on every subscription
// UNSAFE: deferred — UUID computed on every subscription
val mono2 = Mono.fromCallable { UUID.randomUUID().toString() } // new UUID per subscription
// ALSO UNSAFE: deferred inside fromCallable mixed with legitimate deferred work
val mono3 = Mono.fromCallable {
val key = UUID.randomUUID().toString() // new UUID per subscription
computeSomeBlockingHash(key) // blocking work, legitimately deferred
key // re-generated key returned per subscription
}
The fix: compute the key before the transactional() boundary
Compute the UUID (or content-hash key) outside the transactional() block and outside any operator that retryWhen() can re-subscribe. In the eager form, compute it before the reactive pipeline is assembled. In the deferred form, use a Mono.just(key) that carries the pre-computed key into the pipeline.
// BillingService.kt — fixed mode 3: key computed before the transactional() boundary
@Service
class BillingService(
private val transactionalOperator: TransactionalOperator,
private val billingRepository: BillingAttemptRepository,
private val stripeClient: ReactiveStripeService
) {
fun chargeCustomer(customer: Customer): Mono<String> {
// Key computed before the transactional() Mono is assembled.
// This is an eager computation — runs once at the point chargeCustomer() is called.
// The same key value is used on every re-subscription of the transactional() Mono.
val idempotencyKey = contentHashKey(
customer.id,
customer.billingPeriod,
customer.amountCents
)
return transactionalOperator.transactional(
// Mono.just carries the pre-computed key into the pipeline.
// retryWhen() re-subscription does not re-evaluate Mono.just — same key.
Mono.just(idempotencyKey)
.flatMap { key ->
billingRepository.save(
BillingAttempt(
customerId = customer.id,
stripeKey = key,
billingPeriod = customer.billingPeriod
)
).then(stripeClient.charge(key, customer.amountCents, customer.id))
}
).retryWhen(
Retry.backoff(3, Duration.ofSeconds(1))
.filter { it is StripeConnectException }
)
}
private fun contentHashKey(customerId: String, period: String, amountCents: Long): String {
val input = "$customerId|$period|$amountCents"
val digest = MessageDigest.getInstance("SHA-256").digest(
input.toByteArray(Charsets.UTF_8)
)
return digest.joinToString("") { "%02x".format(it) }
}
}
With Mono.just(idempotencyKey), the key is emitted as a constant value on every subscription. retryWhen() re-subscribes the transactional() Mono, a new transaction is opened, but Mono.just(idempotencyKey) emits the same pre-computed key. billingRepository.save() inserts a new DB row (UUID_A again, because BillingAttempt has a new null ID on each re-subscription unless a unique constraint prevents it). stripeClient.charge(UUID_A) sends the same key — Stripe deduplicates — returns the cached result of ch_A. No ch_B.
The DB row duplication (two BillingAttempt rows, both with stripe_key = UUID_A) is still a problem if not guarded. Add a unique constraint on (customer_id, billing_period) on the billing_attempt table to prevent the second INSERT from committing. With the unique constraint, the second save() call inside retryWhen will throw a DataIntegrityViolationException — which is not a retryable exception — terminating the retry with a clear, diagnosable error rather than a silent duplicate charge.
R2DBC-specific retry patterns: coroutine Flow vs. Reactor operators
Developers working in Kotlin with Spring Data R2DBC often choose between Reactor operators (Flux, Mono) and Kotlin coroutine extensions (Flow, awaitSingle()). Both approaches are available; the idempotency key placement rules are the same but the re-subscription mechanisms differ slightly.
Kotlin Flow + R2DBC: Flow.retry{} re-collects the upstream
Spring Data R2DBC provides Kotlin coroutine extensions that return Flow from repository methods. If you use Flow.retry{} (or Flow.retryWhen{}) on a flow that generates a UUID inside its builder block, the retry re-collects the entire upstream flow from the beginning, re-executing the builder block. This is the same semantics as Mode 2 above, but with Kotlin’s Flow API.
// BillingFlow.kt — unsafe: UUID inside flow builder, retry() re-collects upstream
fun billingFlow(customers: List<Customer>): Flow<String> = flow {
for (customer in customers) {
val key = UUID.randomUUID().toString() // ← re-executed on retry re-collection
billingRepository.save(
BillingAttempt(customerId = customer.id, stripeKey = key)
).awaitSingle()
val result = stripeService.charge(key, customer.amountCents, customer.id).awaitSingle()
emit(result)
}
}.retry(3) { it is StripeConnectException } // ← re-collects the entire flow on retry
// BillingFlow.kt — fixed: key computed outside the retried block
fun billingFlow(customers: List<Customer>): Flow<String> = customers.asFlow()
.flatMapConcat { customer ->
val key = contentHashKey(customer.id, customer.billingPeriod, customer.amountCents)
flow {
billingRepository.save(
BillingAttempt(customerId = customer.id, stripeKey = key)
).awaitSingle()
val result = stripeService.charge(key, customer.amountCents, customer.id)
.awaitSingle()
emit(result)
}.retry(3) { it is StripeConnectException }
// retry is inside flatMapConcat — only the failing customer's flow is retried.
// key is from the outer lambda — not re-computed on retry re-collection.
}
See the Kotlin Flow + Exposed + Stripe post for a detailed treatment of Flow cold collection semantics and the three specific Flow API patterns that regenerate keys (flow builder, map + retry, shared mutable state).
The awaitSingle() bridge and transaction propagation
When mixing coroutine extensions with TransactionalOperator, the bridge .awaitSingle() (or .awaitSingleOrNull()) subscribes to the Mono and suspends the coroutine until the result is available. The subscription is made within the current coroutine’s ReactorContext via currentCoroutineContext()[ReactorContext]. If TransactionalOperator.transactional(mono).awaitSingle() is called inside a coroutine that itself is wrapped in a TransactionalOperator (e.g., via @Transactional suspend fun on the outer caller), the inner transactional() call begins a new transaction because the inner Mono’s reactive context is a new subscription context, not the outer caller’s context.
This means transaction propagation with R2DBC and coroutines requires care: @Transactional suspend fun propagates the transaction through the coroutine’s ReactorContext for direct coroutine calls (via CoroutinesTransactionInterceptor), but explicit transactionalOperator.transactional(mono).awaitSingle() calls inside a @Transactional suspend fun may not participate in the outer transaction depending on whether the reactive context is propagated correctly. If in doubt, use @Transactional consistently rather than mixing TransactionalOperator explicit calls with coroutine-transaction propagation.
Comparison table: the three failure modes
| Mode | Re-execution trigger | Re-execution unit | UUID position | Developer misconception | Stripe impact |
|---|---|---|---|---|---|
1. @Transactional suspend fun + @Retryable |
AOP proceed() on retry |
Entire suspend fun body (new coroutine dispatch) |
Local variable inside suspend fun |
“@Retryable might be coroutine-aware like @Transactional” |
UUID_B → ch_B per retry attempt |
2. flatMap { UUID } .retryWhen() outside flatMap |
retryWhen() re-subscription |
Entire customerFlux + flatMap body for all customers |
Local variable inside flatMap lambda |
“retryWhen() retries only the failed element” |
UUID_B per customer → ch_B for all customers processed in first subscription |
3. transactional(Mono.fromCallable{UUID}).retryWhen() |
retryWhen() re-subscription |
Entire transactional() Mono including UUID callable |
Inside Mono.fromCallable inside transactional() |
“retryWhen() retries from after the committed transactional() block” |
UUID_B → ch_B per retry attempt |
| Mode | DB impact | Detection signal | Fix |
|---|---|---|---|
1. @Transactional suspend fun + @Retryable |
Second BillingAttempt row with UUID_B inserted by second transaction |
Two BillingAttempt rows for same customer + period; two Stripe PI IDs in audit log |
Pass key as parameter computed outside @Retryable proxy |
2. flatMap { UUID } .retryWhen() outside flatMap |
New rows for all customers in customerFlux, not just the failed one |
N duplicate rows (N = customers processed before the failure); N ch_B charges | Move retryWhen() inside flatMap; compute key outside retried chain |
3. transactional(Mono.fromCallable{UUID}).retryWhen() |
Second row with UUID_B; unique constraint on (customer_id, billing_period) prevents if added |
Two rows per customer; Stripe duplicate charge for same customer in same period | Compute key before transactional() Mono assembly; use Mono.just(key) |
Testing R2DBC + coroutine retry pipelines for idempotency
Each failure mode requires a test that asserts the Idempotency-Key header is identical across all retry attempts. WireMock with Reactor’s StepVerifier provides the assertion surface.
// BillingServiceTest.kt — Mode 1: assert same Idempotency-Key header across @Retryable retries
@SpringBootTest
@AutoConfigureWireMock(port = 0)
class BillingServiceMode1Test(
@Autowired private val billingService: BillingService,
@Autowired private val wireMockServer: WireMockServer
) {
@Test
fun `@Retryable should use same idempotency key across all retry attempts`() {
// First attempt: Stripe returns a connect exception
wireMockServer.stubFor(
post(urlPathEqualTo("/v1/payment_intents"))
.inScenario("stripe-retry")
.whenScenarioStateIs(Scenario.STARTED)
.willSetStateTo("second-attempt")
.willReturn(aResponse().withStatus(500).withBody("{\"error\":{\"type\":\"api_connection_error\"}}"))
)
// Second attempt: Stripe returns success
wireMockServer.stubFor(
post(urlPathEqualTo("/v1/payment_intents"))
.inScenario("stripe-retry")
.whenScenarioStateIs("second-attempt")
.willReturn(aResponse()
.withStatus(200)
.withHeader("Content-Type", "application/json")
.withBody("""{"id":"pi_test_001","status":"succeeded"}"""))
)
// Compute the expected key — same logic as the fixed BillingService
val expectedKey = contentHashKey("cus_001", "2026-10", 9900L)
billingService.chargeCustomer("cus_001", 9900L, "2026-10", expectedKey)
// Verify both requests used the SAME Idempotency-Key header
val requests = wireMockServer.findAll(postRequestedFor(urlPathEqualTo("/v1/payment_intents")))
assertThat(requests).hasSize(2)
val keys = requests.map { it.getHeader("Idempotency-Key") }.toSet()
assertThat(keys).hasSize(1) // both requests must use the same key
assertThat(keys.first()).isEqualTo(expectedKey)
}
}
// BillingPipelineTest.kt — Mode 2: assert retryWhen inside flatMap retries only failed customer
@SpringBootTest
@AutoConfigureWireMock(port = 0)
class BillingPipelineMode2Test(
@Autowired private val billingPipeline: BillingPipeline,
@Autowired private val wireMockServer: WireMockServer,
@Autowired private val billingRepository: BillingAttemptRepository
) {
@Test
fun `retryWhen inside flatMap should not re-process already-charged customers`() {
val customers = listOf(
Customer("cus_001", 9900L, "2026-10"),
Customer("cus_002", 4900L, "2026-10"),
Customer("cus_003", 14900L, "2026-10") // ← will fail on first attempt
)
// cus_001 and cus_002: always succeed
wireMockServer.stubFor(post(urlPathEqualTo("/v1/payment_intents"))
.withRequestBody(matchingJsonPath("$.customer", equalTo("cus_001")))
.willReturn(aResponse().withStatus(200)
.withBody("""{"id":"pi_001","status":"succeeded"}""")))
wireMockServer.stubFor(post(urlPathEqualTo("/v1/payment_intents"))
.withRequestBody(matchingJsonPath("$.customer", equalTo("cus_002")))
.willReturn(aResponse().withStatus(200)
.withBody("""{"id":"pi_002","status":"succeeded"}""")))
// cus_003: fail on first attempt, succeed on second
wireMockServer.stubFor(post(urlPathEqualTo("/v1/payment_intents"))
.withRequestBody(matchingJsonPath("$.customer", equalTo("cus_003")))
.inScenario("cus003-retry")
.whenScenarioStateIs(Scenario.STARTED)
.willSetStateTo("retry")
.willReturn(aResponse().withStatus(500).withBody("{\"error\":{\"type\":\"api_connection_error\"}}")))
wireMockServer.stubFor(post(urlPathEqualTo("/v1/payment_intents"))
.withRequestBody(matchingJsonPath("$.customer", equalTo("cus_003")))
.inScenario("cus003-retry")
.whenScenarioStateIs("retry")
.willReturn(aResponse().withStatus(200)
.withBody("""{"id":"pi_003","status":"succeeded"}""")))
StepVerifier.create(billingPipeline.runBillingFlux(Flux.fromIterable(customers)))
.expectNextCount(3)
.verifyComplete()
// cus_001 and cus_002 should have exactly 1 Stripe request each (not re-processed)
val req001 = wireMockServer.findAll(postRequestedFor(urlPathEqualTo("/v1/payment_intents"))
.withRequestBody(matchingJsonPath("$.customer", equalTo("cus_001"))))
assertThat(req001).hasSize(1)
val req002 = wireMockServer.findAll(postRequestedFor(urlPathEqualTo("/v1/payment_intents"))
.withRequestBody(matchingJsonPath("$.customer", equalTo("cus_002"))))
assertThat(req002).hasSize(1)
// cus_003 should have exactly 2 Stripe requests (1 fail + 1 retry)
val req003 = wireMockServer.findAll(postRequestedFor(urlPathEqualTo("/v1/payment_intents"))
.withRequestBody(matchingJsonPath("$.customer", equalTo("cus_003"))))
assertThat(req003).hasSize(2)
// All cus_003 requests must use the same Idempotency-Key
val keys003 = req003.map { it.getHeader("Idempotency-Key") }.toSet()
assertThat(keys003).hasSize(1)
// Only 3 BillingAttempt rows should exist in the database (not 6)
val rowCount = billingRepository.count().block()!!
assertThat(rowCount).isEqualTo(3L)
}
}
// BillingServiceTest.kt — Mode 3: assert Mono.just(key) emits same key on re-subscription
class BillingServiceMode3Test {
@Test
fun `Mono_just computed before transactional should emit same key on re-subscription`() {
// Simulates retryWhen re-subscribing the transactional() Mono twice.
// Both subscriptions should receive the same UUID.
val key = contentHashKey("cus_001", "2026-10", 9900L)
val keyMono = Mono.just(key) // eager — same value on every subscription
val keysEmitted = mutableListOf<String>()
repeat(3) {
keyMono.subscribe { keysEmitted.add(it) }
}
// All three subscriptions (simulating 2 retries) must emit the same key
assertThat(keysEmitted).hasSize(3)
assertThat(keysEmitted.toSet()).hasSize(1)
assertThat(keysEmitted.first()).isEqualTo(key)
// Contrast with Mono.fromCallable — new UUID per subscription
val callableMono = Mono.fromCallable { UUID.randomUUID().toString() }
val callableKeys = mutableListOf<String>()
repeat(3) {
callableMono.subscribe { callableKeys.add(it) }
}
assertThat(callableKeys.toSet()).hasSize(3) // three different UUIDs — the bug
}
}
Keybrake’s role in the R2DBC reactive retry surface
The three failure modes above share a structural property: the duplicate charge is created silently. From Stripe’s perspective, each charge is a legitimate, unique request. From the application’s database perspective, the second BillingAttempt row is a legitimate insert. Neither Stripe nor the application database has the information to know that ch_B is a retry of ch_A under a different key. The duplicate is visible only if you correlate the Idempotency-Key header value, the customer ID, the billing period, and the timestamp across multiple Stripe requests — a query that spans Stripe’s API history, not a single-response assertion.
Keybrake sits in the request path between the Spring application and the Stripe API. Every proxied Stripe request is logged with its Idempotency-Key header, the calling agent key, the timestamp, and the parsed charge amount. The Keybrake dashboard surfaces anomalies: multiple Stripe requests within a short time window for the same agent key and customer, with different Idempotency-Key values. This pattern is invisible in the Stripe Dashboard (which shows two independent PaymentIntents) but is directly visible in Keybrake’s per-agent audit log, where the sequence of keys, timestamps, and customers is queryable together.
The per-agent daily spend cap provides an operational backstop. If a flatMap { UUID }.retryWhen() outside the flatMap fires on a 500-customer billing run, the per-agent daily cap will close the key before all 500 customers receive ch_B. The cap fires based on total Stripe amount proxied through the key in the current UTC day. This does not eliminate the duplicate charges that were processed before the cap fired, but it limits the blast radius from what could be 500 duplicate charges to as few as one before the key is blocked — surfacing the anomaly immediately rather than after end-of-month reconciliation.
The Spring Boot 3 virtual threads post describes the same Keybrake audit log pattern for Mode 3 (StructuredTaskScope + outer retry): the signal is multiple different Idempotency-Key values for the same customer in the same billing period, visible in the proxied request log, within seconds of each other. That signal is structurally identical for all three R2DBC modes above.
Summary: the universal placement rule for reactive + coroutine Stripe pipelines
All three failure modes in this post reduce to the same root cause: UUID.randomUUID() (or any key generation that uses local non-deterministic state) is inside a code block that re-executes on retry. The precise mechanism — AOP proceed(), Reactor retryWhen() re-subscription, or Mono.fromCallable re-execution — determines what “inside the re-executed block” means syntactically. But the rule is the same in all cases:
Compute the idempotency key at the outermost scope where it will be used — outside every boundary that can re-execute the code containing the key. For AOP-managed methods, that means outside the proxy boundary (in the caller, passed as a parameter). For reactive pipelines, that means outside the
retryWhen()-retried operator subtree (computed before the pipeline is assembled, carried in asMono.just(key)). ForflatMap-based batch pipelines, that means outside the retried per-element chain (computed in the outerflatMaplambda, withretryWhen()inside the per-element chain).
The content-hash key — derived deterministically from customer ID, billing period, and amount via SHA-256 — satisfies all three cases simultaneously. It can be computed at any time, by any component, and always produces the same output for the same billing event. Stripe deduplicates on it correctly. It is safe to pass across AOP proxy boundaries, reactive context boundaries, and coroutine dispatch boundaries.
See also: Kotlin Coroutines @Transactional + Stripe, Kotlin Flow + Exposed + Stripe, Spring Batch + Stripe, and Spring WebClient reactive retry + Stripe.
Put a spend cap on your agent’s Stripe key
Keybrake proxies your Stripe calls with per-agent spend caps, endpoint allowlists, and a full audit log. Duplicate idempotency keys are visible in the audit log before they become customer support problems.