Spring WebFlux Functional Endpoints and Stripe Integration: How bodyToMono().cache() + Outer retryWhen() Re-invokes the flatMap Lambda, WebClient .headers{} Customizer Regenerates Idempotency Keys on Re-subscription, and Mono.fromCallable() in Per-Item Batch Retry Creates Duplicate Charges

Spring WebFlux supports two programming models: the annotation-based model (@RestController, @GetMapping) and the functional model (RouterFunction, HandlerFunction). The functional model’s idioms for reading request bodies, building WebClient calls, and composing reactive pipelines introduce three Stripe idempotency key failure modes that are structurally different from the modes covered in the general Spring WebFlux post. All three modes here produce the same outcome — UUID.randomUUID() generates UUID_B on a retry attempt, Stripe processes it as a new billing intent, and charge ch_B is committed alongside the already-committed ch_A — but the mechanism in each case is specific to how functional endpoints handle request bodies, how WebClient’s builder API evaluates its configuration lambdas, and how Mono.fromCallable() interacts with per-item retry in a batch handler.

These modes are distinct from those covered in the Spring WebFlux (annotation model) post, which covered Mono.defer() factory re-evaluation, Flux.retryWhen() at the wrong level for batch billing, and @Scheduled returning Mono<Void> being silently discarded. The functional endpoint patterns here — bodyToMono().cache() combined with outer retryWhen(), the WebClient .headers{} customizer lambda, and Mono.fromCallable() in a per-item batch flatMap — are also structurally different from the Spring WebClient + @Transactional post (which covered R2DBC reactive transactions and @Retryable wrapping block()). The audience for this post is teams writing the Spring WebFlux functional style (common in teams that want explicit routing, testable handler functions, and a code-as-configuration approach to request handling) who have added retry for resilience and discovered, or are about to discover, that their retry code generates new idempotency keys on every attempt.

Background: the Spring WebFlux functional endpoint model and why retry interacts differently with it

In the Spring WebFlux functional model, HTTP routes are defined as a RouterFunction<ServerResponse> bean, and the handling logic lives in a HandlerFunction<ServerResponse> — a single-method interface: fun handle(request: ServerRequest): Mono<ServerResponse>. The handler receives a ServerRequest and returns a Mono<ServerResponse>. Reading the request body is done through request.bodyToMono(SomeClass::class.java) or request.bodyToFlux(SomeClass::class.java).

The critical constraint that shapes all three failure modes in this post: request.bodyToMono() creates a cold publisher that reads from the underlying HTTP request body stream. That stream is a network buffer that can only be consumed once. If you apply retryWhen() to a chain that includes bodyToMono(), the retry’s re-subscription attempts to re-read the stream, which is already exhausted. This forces functional endpoint developers to either:

  1. Apply .cache() to the bodyToMono() publisher to convert it to a hot (caching) publisher that replays the already-parsed value on re-subscription, then apply retryWhen() downstream.
  2. Place the retryWhen() inside the flatMap lambda, scoped only to the Stripe call chain, not to the bodyToMono() step.

Option 1 looks correct but contains a trap: .cache() fixes the body-consumed-once problem, but does nothing to protect idempotency keys generated inside the flatMap lambda downstream of the cache. Option 2 is structurally correct for idempotency key stability, but only if the UUID is computed outside the flatMap lambda entirely. If UUID is inside the flatMap lambda but outside the inner retryWhen scope, the inner retry is safe — but only for a single layer of retry. An outer retry (in a calling method, scheduler, or a second retryWhen() operator) will re-invoke the outer flatMap lambda and regenerate UUID.

The three modes below each represent a different way this constraint manifests in functional endpoint code.

Mode 1: bodyToMono().cache() + outer retryWhen() — .cache() fixes the body-consumed-once problem but the flatMap lambda still re-executes — UUID inside the lambda generates UUID_B — ch_B

A developer writing a functional billing endpoint begins with a handler that reads the charge request from the body, calls Stripe, and returns the result. They want to add retry for resilience against transient Stripe 503s. The initial (broken) attempt:

// BillingHandler.kt — unsafe mode 1a: retryWhen() applied directly to bodyToMono() chain
// This fails with "Only one connection receive subscriber allowed" on retry:
// the body stream has been consumed on the first subscription attempt.
class BillingHandlerV1(
    private val webClient: WebClient,
    private val stripeKey: String
) : HandlerFunction<ServerResponse> {

    override fun handle(request: ServerRequest): Mono<ServerResponse> {
        return request.bodyToMono(ChargeRequest::class.java)
            .flatMap { req ->
                val idempotencyKey = UUID.randomUUID().toString()
                callStripe(req, idempotencyKey)
            }
            .retryWhen(Retry.backoff(3, Duration.ofMillis(300))
                .filter { it is StripeException })
            .flatMap { charge -> ServerResponse.ok().bodyValue(charge) }
    }
}

The developer runs this and finds that retry fails with a reactive stream error when Stripe returns a 503: the body buffer was fully consumed on the first attempt, and the retry’s re-subscription to bodyToMono() hits an exhausted stream. The developer correctly diagnoses the problem and applies .cache():

// BillingHandler.kt — unsafe mode 1b: .cache() fixes body-consumed-once,
// but UUID is still inside the flatMap lambda — retryWhen re-invokes the lambda — UUID_B.
class BillingHandlerV1b(
    private val webClient: WebClient,
    private val stripeKey: String
) : HandlerFunction<ServerResponse> {

    override fun handle(request: ServerRequest): Mono<ServerResponse> {
        // .cache() makes the parsed body replay on re-subscription.
        // Developer's reasoning: "Now retryWhen can re-subscribe and get the same ChargeRequest
        // value. The body-consumed-once problem is solved."
        val cachedBody = request.bodyToMono(ChargeRequest::class.java).cache()

        return cachedBody
            .flatMap { req ->
                // Developer's reasoning: "UUID is computed inside the flatMap, which runs once
                // per ChargeRequest value emitted by cachedBody. cachedBody emits exactly one
                // value — the cached ChargeRequest. So UUID runs once."
                // BUG: the flatMap LAMBDA is the unit that retryWhen re-invokes.
                // retryWhen re-subscribes to cachedBody.flatMap{...}.
                // cachedBody replays the same ChargeRequest value (correct).
                // flatMap re-invokes its lambda with that replayed value (BUG source).
                // Every statement inside the lambda body re-executes per invocation.
                // val here is a NEW binding per lambda invocation — UUID.randomUUID()
                // generates UUID_A on invocation 1 and UUID_B on invocation 2 (retry).
                val idempotencyKey = UUID.randomUUID().toString() // UNSAFE

                callStripe(req, idempotencyKey)
            }
            .retryWhen(Retry.backoff(3, Duration.ofMillis(300))
                .filter { it is StripeException })
            .flatMap { charge -> ServerResponse.ok().bodyValue(charge) }
    }

    private fun callStripe(req: ChargeRequest, idempotencyKey: String): Mono<StripeChargeResponse> {
        return webClient.post()
            .uri("https://api.stripe.com/v1/charges")
            .header("Authorization", "Bearer $stripeKey")
            .header("Idempotency-Key", idempotencyKey)
            .body(BodyInserters.fromFormData(
                LinkedMultiValueMap<String, String>().apply {
                    add("amount", req.amountCents.toString())
                    add("currency", "usd")
                    add("customer", req.customerId)
                }
            ))
            .retrieve()
            .bodyToMono(StripeChargeResponse::class.java)
    }
}

The call sequence on a transient Stripe 503:

  1. handle(request) is called by the WebFlux dispatcher. cachedBody is a Mono<ChargeRequest> backed by a caching operator.
  2. The framework subscribes to the returned Mono<ServerResponse>. Reactor evaluates the operator chain bottom-up: subscribes to cachedBody.flatMap{...}, which subscribes to cachedBody.
  3. cachedBody reads the HTTP body stream, parses the ChargeRequest, and emits it. The cache stores this value for future subscribers.
  4. The flatMap lambda is invoked with the ChargeRequest value. Inside: UUID.randomUUID() generates UUID_A. callStripe(req, UUID_A) is called. WebClient sends the POST to Stripe with Idempotency-Key: UUID_A.
  5. Stripe receives the request, processes it, commits charge ch_A to its datastore. Before the 200 response arrives, a network connection drop causes WebClient to throw WebClientResponseException (for a 503 that arrived) — caught by the .filter { it is StripeException } predicate.
  6. retryWhen() fires. It re-subscribes to cachedBody.flatMap{...}. This re-subscription propagates upstream to cachedBody, which replays the already-parsed ChargeRequest from its cache. The flatMap mapper lambda is re-invoked with the same ChargeRequest value.
  7. Re-invocation of the flatMap lambda: the lambda body executes from the top. val idempotencyKey = UUID.randomUUID().toString() — a new val binding, a new random UUID. UUID_B is generated. callStripe(req, UUID_B) is called. WebClient sends the POST to Stripe with Idempotency-Key: UUID_B.
  8. Stripe receives the request with UUID_B. Stripe looks up UUID_B in its idempotency store: no record (ch_A was stored under UUID_A). Stripe processes the request as a new billing intent and commits charge ch_B.
  9. The retry succeeds. The handler emits a ServerResponse with ch_B’s ID. The client sees a successful response. ch_A exists in Stripe but was never recorded in the application’s database. The customer has been charged twice.

The developer’s mental model failure: “I applied .cache() to fix the body-consumed-once problem. Now cachedBody emits the same value on re-subscription. The flatMap processes each emitted value once. Since cachedBody emits exactly one value per subscription, the flatMap runs exactly once per re-subscription. So the whole chain runs once on the first attempt and once on the retry — two runs total, but with the same ChargeRequest input.” The error in this reasoning is conflating “the same input value” with “the same lambda invocation.” cachedBody does replay the same ChargeRequest value. But replaying a value to a flatMap means invoking the lambda again. Each invocation of the lambda body is a fresh activation frame. UUID.randomUUID() is a method call, not a constant — it generates a new value on each call, regardless of what input is passed to the surrounding lambda.

The fix: UUID computed outside the flatMap lambda, before the retry chain

The fix is to compute the idempotency key before the reactive chain and capture it as a closure variable. Because the key is computed at handler invocation time (before any subscription), it does not re-evaluate on retry re-subscription:

// BillingHandler.kt — safe mode 1: content-hash key outside the reactive chain
class BillingHandlerSafe(
    private val webClient: WebClient,
    private val stripeKey: String
) : HandlerFunction<ServerResponse> {

    override fun handle(request: ServerRequest): Mono<ServerResponse> {
        val cachedBody = request.bodyToMono(ChargeRequest::class.java).cache()

        return cachedBody
            .flatMap { req ->
                // Content-hash key: computed from stable billing parameters.
                // billingPeriod is extracted from the request body, not from a random source.
                // This key is the same every time this specific billing intent is retried.
                val idempotencyKey = "charge:${req.billingPeriod}:${req.customerId}:${req.amountCents}"

                callStripe(req, idempotencyKey)
            }
            .retryWhen(Retry.backoff(3, Duration.ofMillis(300))
                .filter { it is StripeException })
            .flatMap { charge -> ServerResponse.ok().bodyValue(charge) }
    }
}

The content-hash key "charge:${req.billingPeriod}:${req.customerId}:${req.amountCents}" is a deterministic function of the billing intent’s identity. The same customer + amount + period combination always produces the same string. When retryWhen() re-invokes the flatMap lambda with the same ChargeRequest value (replayed from .cache()), the key expression evaluates to the same string as attempt 1. Stripe looks up this key, finds ch_A, and returns ch_A’s response without creating ch_B.

Alternative fix: place the retryWhen() inside the flatMap lambda, scoping the retry to the WebClient call only. This approach removes the dependency on .cache() for retry safety (though you may still want .cache() if the handler can be subscribed multiple times for other reasons):

// BillingHandler.kt — alternative safe fix: retryWhen inside flatMap, outside UUID
class BillingHandlerAltSafe(
    private val webClient: WebClient,
    private val stripeKey: String
) : HandlerFunction<ServerResponse> {

    override fun handle(request: ServerRequest): Mono<ServerResponse> {
        return request.bodyToMono(ChargeRequest::class.java)
            .flatMap { req ->
                // UUID is inside the flatMap lambda but outside the retryWhen scope.
                // retryWhen is applied only to the callStripe() Mono, not to bodyToMono().
                // retryWhen re-subscribes to callStripe()'s Mono only — the flatMap lambda
                // is NOT re-invoked. The idempotencyKey binding is captured once.
                val idempotencyKey = "charge:${req.billingPeriod}:${req.customerId}:${req.amountCents}"

                callStripe(req, idempotencyKey)
                    .retryWhen(Retry.backoff(3, Duration.ofMillis(300))
                        .filter { it is StripeException })
            }
            .flatMap { charge -> ServerResponse.ok().bodyValue(charge) }
    }
}

This structure is safe because the retryWhen() is downstream of callStripe() but upstream of the outer flatMap. It re-subscribes only to callStripe()’s internal WebClient chain, not to bodyToMono() and not to the outer flatMap mapper lambda. The idempotencyKey val is bound once per outer flatMap invocation and captured by the callStripe call. All retry attempts reuse the same captured value.

The outer-caller trap: safe inner retry does not protect against outer retry

The alternative safe fix above is correct for retry within the handler. It is not correct if an outer mechanism also retries the handler invocation:

// BillingScheduler.kt — outer retry calls the handler method again
// Each call to billingService.charge() is a fresh method invocation.
// If billingService.charge() uses UUID.randomUUID() at the method scope (not content-hash),
// the UUID regenerates on each outer retry call.
@Scheduled(fixedDelay = 60_000L)
fun runBillingBatch() {
    val failedCustomers = billingRepository.findFailedBillingAttempts()
    failedCustomers.forEach { customer ->
        var outerAttempts = 0
        while (outerAttempts < 3) {
            outerAttempts++
            try {
                billingService.charge(customer.id, customer.amountCents, customer.billingPeriod).block()
                break
            } catch (e: StripeException) {
                if (outerAttempts >= 3) throw e
                Thread.sleep(5000L * outerAttempts)
            }
        }
    }
}

If billingService.charge() uses UUID.randomUUID() at its method scope (instead of a content-hash key), each outer while-loop retry calls billingService.charge() again, which creates a fresh method invocation, which re-evaluates UUID.randomUUID() — UUID_B — ch_B. The inner retryWhen() inside the handler correctly handles transient Stripe network errors. The outer while-loop is an entirely separate retry layer that bypasses the inner retry by re-invoking the entire method. A content-hash key derived from customer.billingPeriod + customer.id + customer.amountCents returns the same string on every call with the same inputs, regardless of how many nested retry layers surround the call site.

Mode 2: WebClient .headers{} customizer lambda — the customizer re-executes on every HTTP request assembly per subscription — retryWhen() re-subscribes — UUID inside the customizer generates UUID_B — ch_B

In Spring WebFlux functional endpoints, teams often extract the Stripe HTTP call into a private helper method or a standalone service that returns a Mono<StripeChargeResponse>. The WebClient API uses a fluent builder pattern where headers, URI parameters, and body are provided as functional interfaces (lambdas) that WebClient invokes when assembling the actual HTTP request. When retryWhen() triggers a retry, WebClient reassembles the entire HTTP request — reinvoking all builder lambdas, including the .headers{} customizer.

// StripeGateway.kt — unsafe mode 2: UUID inside .headers{} customizer lambda
class StripeGateway(
    private val webClient: WebClient,
    private val stripeKey: String
) {

    fun charge(customerId: String, amountCents: Long): Mono<StripeChargeResponse> {
        return webClient.post()
            .uri("https://api.stripe.com/v1/charges")
            .headers { h ->
                // The .headers{} block is a Consumer<HttpHeaders>.
                // WebClient invokes this consumer when assembling the HTTP request.
                // HTTP request assembly happens on every subscription.
                // retryWhen() re-subscribes on Stripe error → WebClient reassembles the request
                // → this lambda is called again → UUID.randomUUID() generates UUID_B.
                h.set("Authorization", "Bearer $stripeKey")
                h.set("Idempotency-Key", UUID.randomUUID().toString()) // UNSAFE
            }
            .body(BodyInserters.fromFormData(
                LinkedMultiValueMap<String, String>().apply {
                    add("amount", amountCents.toString())
                    add("currency", "usd")
                    add("customer", customerId)
                }
            ))
            .retrieve()
            .bodyToMono(StripeChargeResponse::class.java)
            .retryWhen(Retry.backoff(3, Duration.ofMillis(300))
                .filter { it is WebClientResponseException && it.statusCode.is5xxServerError })
    }
}

The failure sequence:

  1. A handler calls stripeGateway.charge(customerId, amountCents). The method builds the WebClient spec and appends .retryWhen(). This returns a Mono<StripeChargeResponse>. The WebClient chain has not executed yet — it is cold.
  2. The framework subscribes to the handler’s returned Mono<ServerResponse>, which eventually subscribes to the Mono<StripeChargeResponse> from charge().
  3. On subscription, WebClient assembles the HTTP request: evaluates the .uri() string (fixed here), invokes the .headers{} customizer lambda. Inside: UUID.randomUUID() generates UUID_A. The request is sent to Stripe with Idempotency-Key: UUID_A.
  4. Stripe processes the request and commits ch_A. Before the 200 response arrives, a 503 is returned. WebClient throws WebClientResponseException.ServiceUnavailable. The retryWhen() filter matches (5xx server error).
  5. retryWhen() fires. It re-subscribes upstream. Reactor re-subscribes to the WebClient request chain from the beginning of the webClient.post()... operator sequence.
  6. WebClient reassembles the HTTP request for the retry attempt: it re-evaluates the .uri() string and re-invokes the .headers{} customizer lambda. Inside: UUID.randomUUID() generates UUID_B — a fresh random value unrelated to UUID_A. The request is sent to Stripe with Idempotency-Key: UUID_B.
  7. Stripe looks up UUID_B: no record. Stripe creates ch_B. The retry succeeds, returning ch_B’s ID. ch_A exists in Stripe, unrecorded in the application database. Duplicate charge.

The developer’s mental model failure: “I set the idempotency key in the .headers{} block, which is part of the request configuration. The request configuration is set once when I build the WebClient spec. The retry retries the same request.” The error: the .headers{} block is not request configuration in the sense of “a value stored in a field and reused.” It is a callback — a Consumer<HttpHeaders> instance — that WebClient calls every time it needs to assemble the headers for an HTTP request. When retryWhen() re-subscribes, WebClient must assemble a new HttpHeaders object for the new HTTP attempt, and it does so by calling the callback again. The callback contains UUID.randomUUID(), which is a method call that returns a new value each time it is evaluated.

The subtler variant: UUID moved outside .headers{} but inside the helper method — helper method called from a flatMap lambda subject to outer retry

A developer reading about this failure mode “fixes” it by moving the UUID computation outside the .headers{} block:

// StripeGateway.kt — partially fixed: UUID outside .headers{} but inside the helper method
class StripeGatewayPartialFix(
    private val webClient: WebClient,
    private val stripeKey: String
) {

    fun charge(customerId: String, amountCents: Long): Mono<StripeChargeResponse> {
        // UUID is now outside the .headers{} lambda — stable within this method call.
        // retryWhen() applied inside this method re-subscribes to the WebClient chain.
        // The .headers{} lambda now captures idempotencyKey from the method scope.
        // On re-subscription, WebClient re-invokes .headers{} — but idempotencyKey
        // is a captured val from the method scope (evaluated once when charge() was called).
        // So .headers{} correctly uses UUID_A on all retry attempts.
        // THIS FIX IS CORRECT for the inner retryWhen() inside charge().
        val idempotencyKey = UUID.randomUUID().toString() // outside .headers{} — CORRECT for inner retry

        return webClient.post()
            .uri("https://api.stripe.com/v1/charges")
            .headers { h ->
                h.set("Authorization", "Bearer $stripeKey")
                h.set("Idempotency-Key", idempotencyKey) // closure capture — stable for inner retry
            }
            .body(BodyInserters.fromFormData(
                LinkedMultiValueMap<String, String>().apply {
                    add("amount", amountCents.toString())
                    add("currency", "usd")
                    add("customer", customerId)
                }
            ))
            .retrieve()
            .bodyToMono(StripeChargeResponse::class.java)
            .retryWhen(Retry.backoff(3, Duration.ofMillis(300))
                .filter { it is WebClientResponseException && it.statusCode.is5xxServerError })
    }
}

This partial fix is correct for the inner retryWhen(). The UUID is computed once when charge() is called, captured as a closure by the .headers{} lambda, and reused on all retry re-subscriptions within that single charge() call. However, this fix creates a different vulnerability when charge() is called from a flatMap lambda that is itself subject to an outer retryWhen():

// BillingHandler.kt — outer retryWhen() re-invokes flatMap lambda — flatMap calls charge()
// — charge() re-invokes UUID.randomUUID() at its method scope — UUID_B — ch_B
class BillingHandlerOuterRetry(
    private val stripeGateway: StripeGateway,
) : HandlerFunction<ServerResponse> {

    override fun handle(request: ServerRequest): Mono<ServerResponse> {
        val cachedBody = request.bodyToMono(ChargeRequest::class.java).cache()

        return cachedBody
            .flatMap { req ->
                // stripeGateway.charge() computes UUID.randomUUID() at method scope.
                // This is called once per flatMap lambda invocation.
                // The outer retryWhen() re-invokes this flatMap lambda on retry.
                // Each flatMap re-invocation calls stripeGateway.charge() again.
                // charge() computes UUID.randomUUID() again — UUID_B — ch_B.
                stripeGateway.charge(req.customerId, req.amountCents) // generates UUID per call
            }
            .retryWhen(Retry.backoff(2, Duration.ofMillis(500))
                .filter { it is WebClientResponseException && it.statusCode.is5xxServerError })
            .flatMap { charge -> ServerResponse.ok().bodyValue(charge) }
    }
}

The outer retryWhen() (applied to the outer cachedBody.flatMap{...} chain) re-invokes the outer flatMap lambda on retry. The outer lambda calls stripeGateway.charge(req.customerId, req.amountCents). Each call to charge() is a fresh method invocation — a fresh activation frame — in which val idempotencyKey = UUID.randomUUID().toString() re-evaluates. UUID_B. The inner retryWhen() inside charge() correctly retries HTTP-level failures within a single charge invocation, but the outer retry at the handler level creates a second, separate charge attempt with a new UUID.

The fix: content-hash key as a parameter to the helper method

// StripeGateway.kt — safe: idempotency key as a parameter, computed by the caller
class StripeGatewaySafe(
    private val webClient: WebClient,
    private val stripeKey: String
) {

    fun charge(customerId: String, amountCents: Long, idempotencyKey: String): Mono<StripeChargeResponse> {
        return webClient.post()
            .uri("https://api.stripe.com/v1/charges")
            .headers { h ->
                h.set("Authorization", "Bearer $stripeKey")
                h.set("Idempotency-Key", idempotencyKey) // caller-provided — stable
            }
            .body(BodyInserters.fromFormData(
                LinkedMultiValueMap<String, String>().apply {
                    add("amount", amountCents.toString())
                    add("currency", "usd")
                    add("customer", customerId)
                }
            ))
            .retrieve()
            .bodyToMono(StripeChargeResponse::class.java)
            .retryWhen(Retry.backoff(3, Duration.ofMillis(300))
                .filter { it is WebClientResponseException && it.statusCode.is5xxServerError })
    }
}

// BillingHandler.kt — safe: content-hash key computed before the reactive chain
class BillingHandlerSafeV2(
    private val stripeGateway: StripeGatewaySafe,
) : HandlerFunction<ServerResponse> {

    override fun handle(request: ServerRequest): Mono<ServerResponse> {
        val cachedBody = request.bodyToMono(ChargeRequest::class.java).cache()

        return cachedBody
            .flatMap { req ->
                // Content-hash key: deterministic function of billing intent identity.
                // Computed inside the flatMap lambda, but from stable req fields — not UUID.randomUUID().
                // The outer retryWhen() re-invokes this flatMap lambda with the same req value.
                // The key expression evaluates to the same string on every invocation.
                // stripeGateway.charge() uses the caller-provided key — no UUID generation inside.
                val idempotencyKey = "charge:${req.billingPeriod}:${req.customerId}:${req.amountCents}"
                stripeGateway.charge(req.customerId, req.amountCents, idempotencyKey)
            }
            .retryWhen(Retry.backoff(2, Duration.ofMillis(500))
                .filter { it is WebClientResponseException && it.statusCode.is5xxServerError })
            .flatMap { charge -> ServerResponse.ok().bodyValue(charge) }
    }
}

With this structure, the content-hash key is computed inside the flatMap lambda, but it is computed from req.billingPeriod, req.customerId, and req.amountCents — all stable fields that come from the cached ChargeRequest body. When the outer retryWhen() re-invokes the flatMap lambda with the same ChargeRequest value (replayed from .cache()), the key expression evaluates to the same string. stripeGateway.charge() receives the same key. The .headers{} customizer inside charge() captures it as a closure and sends it on every retry. Stripe sees the same key on every attempt and deduplicates correctly.

Mode 3: Mono.fromCallable() inside per-item batch flatMap + per-item retryWhen() — fromCallable() re-executes its callable on every subscription — UUID inside the callable generates UUID_B — ch_B

A common Spring WebFlux functional endpoint pattern for batch billing is a handler that reads a list of customers from the request body, maps each to a Mono via flatMap, and collects the results. Per-item retry is added to each customer’s Mono using retryWhen() inside the flatMap lambda. Teams frequently use Mono.fromCallable() to wrap a blocking or lambda-based Stripe operation into a reactive type. Mono.fromCallable() creates a cold publisher that executes its callable on each subscription:

// BatchBillingHandler.kt — unsafe mode 3: UUID inside Mono.fromCallable() with per-item retryWhen()
class BatchBillingHandler(
    private val webClient: WebClient,
    private val stripeKey: String
) : HandlerFunction<ServerResponse> {

    override fun handle(request: ServerRequest): Mono<ServerResponse> {
        return request.bodyToMono(BatchChargeRequest::class.java)
            .flatMapMany { batch ->
                Flux.fromIterable(batch.customers)
                    .flatMap { customer -> billCustomer(customer) }
            }
            .collectList()
            .flatMap { results -> ServerResponse.ok().bodyValue(BatchResult(results)) }
    }

    private fun billCustomer(customer: CustomerCharge): Mono<ChargeResult> {
        // Mono.fromCallable() wraps a lambda into a cold Mono.
        // The callable body executes on every subscription to this Mono.
        // Per-item retryWhen() below re-subscribes to this Mono on Stripe error.
        // On re-subscription, fromCallable() calls its callable again.
        // UUID.randomUUID() inside the callable generates UUID_B on the retry subscription.
        return Mono.fromCallable {
            val idempotencyKey = UUID.randomUUID().toString() // UNSAFE — inside fromCallable

            ChargeSpec(
                customerId = customer.id,
                amountCents = customer.amountCents,
                idempotencyKey = idempotencyKey
            )
        }
        .flatMap { spec -> callStripe(spec) }
        .retryWhen(Retry.backoff(2, Duration.ofMillis(200))
            .filter { it is WebClientResponseException && it.statusCode.is5xxServerError })
        .map { response -> ChargeResult(customer.id, response.id, "succeeded") }
        .onErrorReturn(ChargeResult(customer.id, null, "failed"))
    }

    private fun callStripe(spec: ChargeSpec): Mono<StripeChargeResponse> {
        return webClient.post()
            .uri("https://api.stripe.com/v1/charges")
            .header("Authorization", "Bearer $stripeKey")
            .header("Idempotency-Key", spec.idempotencyKey)
            .body(BodyInserters.fromFormData(
                LinkedMultiValueMap<String, String>().apply {
                    add("amount", spec.amountCents.toString())
                    add("currency", "usd")
                    add("customer", spec.customerId)
                }
            ))
            .retrieve()
            .bodyToMono(StripeChargeResponse::class.java)
    }
}

The per-item retryWhen() is correctly scoped — it is inside billCustomer(), applied to the chain starting from Mono.fromCallable(). This means only the individual customer’s billing is retried, not the entire batch. This is the correct structure for per-item retry. But Mono.fromCallable() creates a fundamentally cold publisher: its callable is the factory that produces the value on each new subscription. When retryWhen() re-subscribes, the callable re-executes, including UUID.randomUUID().

The failure sequence for a single customer whose first attempt triggers a Stripe 503:

  1. The batch handler processes customer: CustomerCharge(id="cus_001", amountCents=9900).
  2. billCustomer(customer) is called. The returned Mono<ChargeResult> is cold — nothing has executed yet.
  3. The flatMap in the batch handler subscribes to billCustomer()’s Mono. Subscription propagates upstream to Mono.fromCallable(). The callable executes: UUID.randomUUID() generates UUID_A. ChargeSpec(customerId="cus_001", amountCents=9900, idempotencyKey=UUID_A) is emitted.
  4. The downstream flatMap { spec -> callStripe(spec) } receives the spec and calls callStripe() with UUID_A. WebClient sends the POST to Stripe with Idempotency-Key: UUID_A.
  5. Stripe processes the request, commits ch_A. A 503 response arrives. WebClient throws WebClientResponseException.ServiceUnavailable. The retryWhen() filter matches.
  6. retryWhen() re-subscribes upstream. Re-subscription propagates to Mono.fromCallable(). The callable executes again: UUID.randomUUID() generates UUID_B. ChargeSpec(..., idempotencyKey=UUID_B) is emitted.
  7. The downstream flatMap calls callStripe() with UUID_B. WebClient sends the POST to Stripe with Idempotency-Key: UUID_B. Stripe looks up UUID_B: no record. Creates ch_B. Customer charged twice.

The developer’s mental model failure: “I used Mono.fromCallable() to wrap the idempotency key computation and spec building into a reactive step. The key is computed inside the callable, which is part of the Mono pipeline. I then applied retryWhen() to the chain. But the retryWhen() is applied after the fromCallable step in the chain declaration order — it appears downstream in the code. I thought that downstream operators only affect how the result is handled, not whether the upstream re-executes.” The error: in Reactor, a downstream retryWhen() retries by re-subscribing to the upstream. Re-subscribing to an upstream Mono.fromCallable() re-executes its callable. The callable is the upstream. The chain declaration order (fromCallable → flatMap → retryWhen) means retryWhen is declared last, but at runtime it re-subscribes to fromCallable first.

The fix: UUID outside Mono.fromCallable() — or use Mono.just(stableKey) as the source

There are two equivalent fixes. First, compute the content-hash key before the Mono.fromCallable() call and pass it in:

// BatchBillingHandler.kt — safe mode 3a: UUID outside fromCallable
private fun billCustomer(customer: CustomerCharge): Mono<ChargeResult> {
    // Content-hash key computed before the reactive chain.
    // billingPeriod is a stable field on CustomerCharge — passed from the batch request body.
    // This key is deterministic: same customer + period + amount = same key, every time.
    val idempotencyKey = "charge:${customer.billingPeriod}:${customer.id}:${customer.amountCents}"

    val spec = ChargeSpec(
        customerId = customer.id,
        amountCents = customer.amountCents,
        idempotencyKey = idempotencyKey
    )

    // Mono.just() emits the same already-constructed spec on every subscription.
    // retryWhen() re-subscribes to Mono.just(spec) — the same spec value is replayed.
    // No re-computation of idempotencyKey occurs.
    return Mono.just(spec)
        .flatMap { s -> callStripe(s) }
        .retryWhen(Retry.backoff(2, Duration.ofMillis(200))
            .filter { it is WebClientResponseException && it.statusCode.is5xxServerError })
        .map { response -> ChargeResult(customer.id, response.id, "succeeded") }
        .onErrorReturn(ChargeResult(customer.id, null, "failed"))
}

Second, use Mono.fromCallable() for the Stripe HTTP call itself (if wrapping a blocking SDK call), but compute the key before the fromCallable boundary:

// BatchBillingHandler.kt — safe mode 3b: fromCallable wraps only the HTTP call, key is outside
private fun billCustomer(customer: CustomerCharge): Mono<ChargeResult> {
    val idempotencyKey = "charge:${customer.billingPeriod}:${customer.id}:${customer.amountCents}"

    // If using Stripe's blocking Java SDK (Charge.create()) rather than WebClient,
    // wrap the blocking call in fromCallable. But compute the key BEFORE this boundary.
    // fromCallable re-executes on re-subscription — but idempotencyKey is a stable
    // captured val computed before the lambda boundary. Each re-execution of the callable
    // uses the same captured idempotencyKey value.
    return Mono.fromCallable {
        val params = ChargeCreateParams.builder()
            .setAmount(customer.amountCents)
            .setCurrency("usd")
            .setCustomer(customer.id)
            .build()
        Charge.create(params, RequestOptions.builder()
            .setIdempotencyKey(idempotencyKey) // stable captured val — safe
            .setApiKey(stripeKey)
            .build())
    }
    .subscribeOn(Schedulers.boundedElastic()) // blocking SDK call needs non-event-loop thread
    .retryWhen(Retry.backoff(2, Duration.ofMillis(200))
        .filter { it is StripeException && (it as? ApiException)?.code?.let { c ->
            c == "rate_limit" || c == "api_connection_error"
        } == true })
    .map { charge -> ChargeResult(customer.id, charge.id, "succeeded") }
    .onErrorReturn(ChargeResult(customer.id, null, "failed"))
}

In this second variant, Mono.fromCallable() re-executes its callable on each retry re-subscription (as before), but idempotencyKey is a captured val from the outer method scope — computed once when billCustomer() is called, not inside the callable. The callable closes over the already-computed key. Each re-execution reads the same captured value.

The outer-Flux-retry blast radius: per-item retry replaced or supplemented with outer Flux retry

A subtler variant appears when the per-item retryWhen() inside billCustomer() is removed (or never added) and an outer Flux.retryWhen() is applied to the batch pipeline:

// BatchBillingHandler.kt — unsafe mode 3 variant: outer Flux retryWhen()
// Per-item retryWhen removed; outer retryWhen added for "simplicity"
class BatchBillingHandlerOuterRetry(
    private val webClient: WebClient,
    private val stripeKey: String
) : HandlerFunction<ServerResponse> {

    override fun handle(request: ServerRequest): Mono<ServerResponse> {
        return request.bodyToMono(BatchChargeRequest::class.java)
            .flatMapMany { batch ->
                Flux.fromIterable(batch.customers)
                    .flatMap { customer ->
                        // UUID is computed inside the outer flatMap lambda — once per customer
                        // for the initial subscription. But if the outer Flux retryWhen() fires,
                        // the entire Flux re-subscribes — Flux.fromIterable() re-emits all customers
                        // from the beginning — this flatMap lambda re-executes for EVERY customer,
                        // including those already successfully charged on the previous subscription.
                        val idempotencyKey = UUID.randomUUID().toString() // UNSAFE with outer retry

                        callStripe(
                            customerId = customer.id,
                            amountCents = customer.amountCents,
                            idempotencyKey = idempotencyKey
                        )
                        .map { response -> ChargeResult(customer.id, response.id, "succeeded") }
                    }
            }
            .retryWhen(Retry.backoff(2, Duration.ofMillis(500))
                .filter { it is WebClientResponseException && it.statusCode.is5xxServerError })
            // Outer retryWhen on the Flux: single failure from any customer restarts the entire batch.
            .collectList()
            .flatMap { results -> ServerResponse.ok().bodyValue(BatchResult(results)) }
    }
}

With outer Flux.retryWhen(): when customer 12 out of 50 returns a Stripe 503, retryWhen() fires. It re-subscribes to the Flux.fromIterable(batch.customers) upstream. The entire Flux restarts from the beginning. Flux.fromIterable() re-emits all 50 customers from position 0. The flatMap lambda re-executes for all 50 customers — including customers 1 through 11, whose charges already succeeded in Stripe on the first subscription. Each flatMap re-execution calls UUID.randomUUID(), generating UUID_B through UUID_K for the already-charged customers. Stripe sees UUID_B for customer 1, no record, creates ch_1B alongside already-committed ch_1A. Customers 1 through 11 each get charged twice. Customer 12’s original charge ch_12A also exists (Stripe processed it before the 503 response); the retry’s new UUID creates ch_12B. Total: 12 or more duplicate charges from a single transient error in customer 12’s billing. The blast radius is “number of customers processed before the failure.”

The fix for the outer-retry variant is to restructure: remove the outer Flux.retryWhen() and add per-item retryWhen() inside the flatMap lambda, after computing a content-hash key for that customer:

// Safe: per-item retry with content-hash keys — outer Flux never retries
Flux.fromIterable(batch.customers)
    .flatMap { customer ->
        val idempotencyKey =
            "charge:${batch.billingPeriod}:${customer.id}:${customer.amountCents}"
        callStripe(customer.id, customer.amountCents, idempotencyKey)
            .retryWhen(Retry.backoff(2, Duration.ofMillis(200)) // per-item retry
                .filter { it is WebClientResponseException && it.statusCode.is5xxServerError })
            .map { response -> ChargeResult(customer.id, response.id, "succeeded") }
            .onErrorReturn(ChargeResult(customer.id, null, "failed"))
    }
    // No outer .retryWhen() on the Flux — per-item retry handles all transient errors
    .collectList()

Each customer’s per-item retryWhen() retries only that customer’s billing. The outer Flux never retries — it sees only success signals (either a successful charge or a ChargeResult("failed") from .onErrorReturn()). Content-hash keys ensure that a customer whose first attempt committed ch_A gets the same key on the per-item retry, causing Stripe to return ch_A without creating ch_B.

Summary: the three functional-endpoint-specific failure modes compared

Mode Mechanism that re-evaluates UUID Why the developer misses it Fix
1 Outer retryWhen() re-subscribes to cachedBody.flatMap{...} — flatMap lambda re-invoked with the replayed body value — UUID inside lambda generates UUID_B “.cache() fixed the body-consumed-once problem; I thought the flatMap runs once per value, not once per subscription” Content-hash key computed from req fields, not UUID.randomUUID(), inside the flatMap lambda; or place retryWhen() inside the flatMap, scoped to the WebClient call only
2 WebClient .headers{} customizer lambda re-invoked on every HTTP request assembly — retryWhen() re-subscribes — WebClient reassembles — UUID inside customizer generates UUID_B “The headers block is request configuration, set once when I build the spec; retry retries the same request” UUID as a val outside .headers{}, captured by the customizer; or pass idempotency key as a parameter to the helper method, computed by the caller from stable billing intent fields
3 Mono.fromCallable() callable re-executes on every subscription — per-item retryWhen() re-subscribes — callable re-evaluates UUID inside it — UUID_B; OR outer Flux.retryWhen() restarts entire batch — flatMap lambda re-executes for already-charged customers — blast radius = already-processed count “fromCallable() runs before retryWhen() in the code; I thought downstream operators can’t cause upstream re-execution” / “One outer retry is simpler than per-item retry” Content-hash key computed before fromCallable() boundary, as a captured val; per-item retryWhen() inside flatMap, not outer Flux.retryWhen()

WebTestClient + WireMock test patterns for all three modes

All three modes are detectable with the same test structure: WireMock stubs Stripe to return a 503 on the first attempt and a 200 on the second, the test captures all requests sent to Stripe, and asserts that the Idempotency-Key header is identical across all attempts. An unsafe implementation fails the equality assertion. A safe implementation passes.

Mode 1: bodyToMono().cache() + outer retryWhen()

// BillingHandlerTest.kt
@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT)
class BillingHandlerTest {

    private val wireMock = WireMockServer(wireMockConfig().dynamicPort())

    @BeforeEach fun start() { wireMock.start() }
    @AfterEach fun stop() { wireMock.stop() }

    @Test
    fun `outer retryWhen sends same idempotency key on all attempts`() {
        // First Stripe call: 503
        wireMock.stubFor(post(urlEqualTo("/v1/charges"))
            .inScenario("charge-retry")
            .whenScenarioStateIs(Scenario.STARTED)
            .willReturn(serverError())
            .willSetStateTo("retried"))

        // Second Stripe call: 200
        wireMock.stubFor(post(urlEqualTo("/v1/charges"))
            .inScenario("charge-retry")
            .whenScenarioStateIs("retried")
            .willReturn(okJson("""{"id":"ch_ok","object":"charge","status":"succeeded"}""")))

        // Build WebTestClient pointing at the WebFlux app.
        // The app's StripeGateway WebClient should be configured to point at wireMock.port.
        val client = WebTestClient
            .bindToRouterFunction(billingRouter(BillingHandlerSafe(wireMockWebClient(), "sk_test")))
            .build()

        client.post().uri("/api/charge")
            .contentType(MediaType.APPLICATION_JSON)
            .bodyValue("""{"customerId":"cus_001","amountCents":9900,"billingPeriod":"2026-10"}""")
            .exchange()
            .expectStatus().isOk

        // Assert: Stripe received exactly 2 requests (1 failed + 1 success)
        val requests = wireMock.findAll(postRequestedFor(urlEqualTo("/v1/charges")))
        assertThat(requests).hasSize(2)

        // Assert: same Idempotency-Key on both attempts
        val keys = requests.map { it.getHeader("Idempotency-Key") }
        assertThat(keys[0]).isEqualTo(keys[1])
        assertThat(keys[0]).isEqualTo("charge:2026-10:cus_001:9900")
    }
}

Mode 2: WebClient .headers{} customizer

// StripeGatewayTest.kt
class StripeGatewayTest {

    private val wireMock = WireMockServer(wireMockConfig().dynamicPort())

    @BeforeEach fun start() { wireMock.start() }
    @AfterEach fun stop() { wireMock.stop() }

    @Test
    fun `charge sends same idempotency key on retry after 503`() {
        wireMock.stubFor(post(urlEqualTo("/v1/charges"))
            .inScenario("gateway-retry")
            .whenScenarioStateIs(Scenario.STARTED)
            .willReturn(serverError())
            .willSetStateTo("retried"))

        wireMock.stubFor(post(urlEqualTo("/v1/charges"))
            .inScenario("gateway-retry")
            .whenScenarioStateIs("retried")
            .willReturn(okJson("""{"id":"ch_ok","object":"charge","status":"succeeded"}""")))

        val webClient = WebClient.builder()
            .baseUrl("http://localhost:${wireMock.port}")
            .build()

        val gateway = StripeGatewaySafe(webClient, "sk_test")

        StepVerifier.create(
            gateway.charge("cus_001", 9900L, "charge:2026-10:cus_001:9900")
        )
        .expectNextMatches { it.id == "ch_ok" }
        .verifyComplete()

        val requests = wireMock.findAll(postRequestedFor(urlEqualTo("/v1/charges")))
        assertThat(requests).hasSize(2)

        val keys = requests.map { it.getHeader("Idempotency-Key") }
        assertThat(keys[0]).isEqualTo(keys[1])
        assertThat(keys[0]).isEqualTo("charge:2026-10:cus_001:9900")
    }
}

Mode 3: Mono.fromCallable() per-item batch retry

// BatchBillingHandlerTest.kt
class BatchBillingHandlerTest {

    private val wireMock = WireMockServer(wireMockConfig().dynamicPort())

    @BeforeEach fun start() { wireMock.start() }
    @AfterEach fun stop() { wireMock.stop() }

    @Test
    fun `per-item retry sends same idempotency key for the retried customer`() {
        // Customer 1: Stripe 503 on first attempt, 200 on second
        wireMock.stubFor(post(urlEqualTo("/v1/charges"))
            .withRequestBody(containing("cus_001"))
            .inScenario("batch-c1")
            .whenScenarioStateIs(Scenario.STARTED)
            .willReturn(serverError())
            .willSetStateTo("retried"))

        wireMock.stubFor(post(urlEqualTo("/v1/charges"))
            .withRequestBody(containing("cus_001"))
            .inScenario("batch-c1")
            .whenScenarioStateIs("retried")
            .willReturn(okJson("""{"id":"ch_001","object":"charge","status":"succeeded"}""")))

        // Customer 2: no failure
        wireMock.stubFor(post(urlEqualTo("/v1/charges"))
            .withRequestBody(containing("cus_002"))
            .willReturn(okJson("""{"id":"ch_002","object":"charge","status":"succeeded"}""")))

        val webClient = WebClient.builder()
            .baseUrl("http://localhost:${wireMock.port}")
            .build()

        val handler = BatchBillingHandlerSafe(webClient, "sk_test")
        val router = RouterFunctions.route(
            RequestPredicates.POST("/api/billing/batch"), handler
        )

        WebTestClient.bindToRouterFunction(router).build()
            .post().uri("/api/billing/batch")
            .contentType(MediaType.APPLICATION_JSON)
            .bodyValue("""
                {
                  "billingPeriod": "2026-10",
                  "customers": [
                    {"id": "cus_001", "amountCents": 9900},
                    {"id": "cus_002", "amountCents": 4900}
                  ]
                }
            """.trimIndent())
            .exchange()
            .expectStatus().isOk

        // cus_001 was retried once — should have 2 requests but same idempotency key
        val c1Requests = wireMock.findAll(
            postRequestedFor(urlEqualTo("/v1/charges")).withRequestBody(containing("cus_001"))
        )
        assertThat(c1Requests).hasSize(2)
        val c1Keys = c1Requests.map { it.getHeader("Idempotency-Key") }
        assertThat(c1Keys[0]).isEqualTo(c1Keys[1])
        assertThat(c1Keys[0]).isEqualTo("charge:2026-10:cus_001:9900")

        // cus_002 was not retried — one request, one key
        val c2Requests = wireMock.findAll(
            postRequestedFor(urlEqualTo("/v1/charges")).withRequestBody(containing("cus_002"))
        )
        assertThat(c2Requests).hasSize(1)
        assertThat(c2Requests[0].getHeader("Idempotency-Key"))
            .isEqualTo("charge:2026-10:cus_002:4900")
    }
}

The assertion pattern that catches the unsafe implementation is assertThat(keys[0]).isEqualTo(keys[1]). When UUID.randomUUID() is inside the retry scope, the two requests have different keys and this assertion fails with a message showing two distinct UUIDs. When the content-hash key is used, both keys are the same deterministic string, and the assertion passes. The test is not checking business logic or happy-path behavior — it is specifically checking the idempotency key stability contract, which is the contract that protects against duplicate charges at the Stripe billing layer.

The relationship between the functional model’s modes and the annotation model’s modes

The three modes in this post are specific to the functional endpoint programming model. They do not appear (in the same form) in annotated controllers:

The unified rule

Across all three modes in this post — and across all other retry patterns covered in this blog series — the same principle applies:

A retry mechanism re-executes a unit of work. Any expression inside that unit re-executes on each retry attempt. UUID.randomUUID() inside that unit generates a new UUID per attempt — UUID_B — which Stripe processes as a new billing intent, creating ch_B alongside any already-committed ch_A.

In the functional endpoint model, the “unit of work” is identified by the retry operator and the cold publisher semantics:

These patterns are related to the modes covered in the general Spring WebFlux post (annotation model), the Spring WebClient + @Transactional post, the Kotlin Coroutines post, and the Ktor + Exposed post. The surface changes (functional endpoint idioms, WebClient builder API, Mono.fromCallable()) but the root cause is the same: the developer’s mental model has UUID generation as a one-time computation, and the framework’s retry model makes it a per-attempt computation by re-invoking the scope it lives in.

Keybrake: put the brakes on your agent’s keys

A scoped API-key proxy for Stripe, Twilio, and Resend — with per-vendor spend caps, endpoint allowlists, per-call audit log, and one-click revoke. Built for teams running autonomous agents against production SaaS APIs.