Micronaut ReactorHttpClient Retry and Stripe Integration: How Mono.create() Sink Factory Re-subscription, Flux.flatMap() Batch Mapper Re-execution, and Mono.defer() Inside retryWhen(Retry.fixedDelay().filter()) Generate New Idempotency Keys

When a Micronaut application uses ReactorHttpClient to call the Stripe API and adds reactive retry logic for resilience, three structurally distinct mechanisms each silently generate a new idempotency key on every retry attempt — and Stripe creates a second charge. All three cases share the same root cause: UUID.randomUUID() is placed inside a scope that re-executes per Reactor subscription, but the Reactor programming model — Mono.create()’s cold sink factory, Flux.fromIterable().flatMap()’s per-item mapper re-execution on batch re-subscription, and Mono.defer()’s factory-per-subscription semantics inside a filtered retryWhen() — each creates a non-obvious per-subscription re-execution boundary that developers routinely place UUID generation inside.

These failure modes are structurally distinct from those covered in the earlier post on Micronaut HTTP client basic patterns (which addressed @ClientFilter per-retry invocation, Mono.fromCallable() with .retry(N), and @Retryable catching database exceptions after a successful Stripe charge), the RxJava3 retry post (Single.fromCallable() re-subscription, flatMap mapper with outer Single.retry(), and @Retryable on Single<T> methods), and the Micronaut Data reactive @Transactional post (Mono.fromCallable() under retryWhen() in reactive transactions, @Retryable proceed() re-invocation on Mono<T>, and onErrorResume() recovery lambda). The three modes here focus on Reactor APIs used specifically with Micronaut’s ReactorHttpClient: the Mono.create() sink programming model for wrapping callback-based async code, the batch Flux pipeline and what happens when the outer Flux retries, and the Mono.defer() + custom-predicate retryWhen() combination where the developer’s mental model of “deferred = per-attempt” is correct about execution timing but wrong about what Stripe’s idempotency system requires.

Background: Micronaut ReactorHttpClient and how Reactor’s cold-publisher contract applies to retry

Micronaut’s ReactorHttpClient is the Project Reactor-compatible variant of Micronaut’s HTTP client. Unlike the blocking HttpClient, ReactorHttpClient returns Mono<HttpResponse<T>> and Flux<HttpResponse<T>>, integrating directly with the Reactor pipeline without any bridge or adapter. In Micronaut applications that already use Project Reactor — particularly those also using Micronaut Data with R2DBC or reactive messaging — it is the natural choice.

Reactor’s core contract is that all publishers returned by Mono and Flux factory methods are cold: nothing happens until a subscriber subscribes, and each new subscription starts a fresh execution. Mono.create(), Mono.defer(), and Mono.fromCallable() are all cold. The distinction between them is how they describe the work to be done per subscription: fromCallable() takes a blocking Callable<T> (runs on a scheduler, delivers the return value); defer() takes a Supplier<Mono<T>> (runs inline, produces another Mono per subscription); create() takes a Consumer<MonoSink<T>> (runs inline, delivers the value via sink.success() or sink.error() callbacks). All three re-execute their factory argument on every new subscription.

Reactor’s retry operators — .retry(N), .retryWhen(Retry) — work by re-subscribing to the upstream Mono or Flux on failure. This is not application-level retry in the sense of “call this service method again”; it is reactive-stream-level retry in the sense of “subscribe to this publisher again from the beginning.” From Reactor’s perspective, a re-subscription is semantically identical to a first subscription. Whatever the publisher does on subscription, it does again.

This distinction matters for Stripe idempotency keys. When UUID.randomUUID() is placed inside a cold-publisher factory that re-executes per subscription, each Reactor retry re-subscription re-evaluates the UUID. Stripe sees a new key and treats the request as a new billing intent — creating a new charge even if the original charge from the first subscription was already committed.

Mode 1: UUID inside Mono.create() sink consumer — wrapping Stripe’s CompletableFuture callback — retryWhen() re-subscribes — sink consumer re-runs — UUID_B — ch_B

The first failure mode arises when a developer uses Mono.create() to wrap Stripe’s asynchronous Java SDK and places UUID.randomUUID() inside the sink consumer. This is the idiomatic Reactor pattern for wrapping callback-based asynchronous APIs — code that delivers its result not via a return value but via a callback function (CompletableFuture.whenComplete(), listener interfaces, completion handlers). The Stripe Java SDK’s Charge.createAsync() returns a CompletableFuture<Charge>, making Mono.create() the natural wrapping choice.

The developer’s reasoning is often: “I’m building a fresh request per subscription — that includes a fresh UUID. Each attempt is independent.” This reasoning would be sound if each Reactor subscription corresponded to a distinct billing intent that Stripe should treat as new. But that is not how Stripe’s idempotency system works: Stripe’s model is that the same billing intent should carry the same idempotency key across all HTTP-level retry attempts, so that Stripe can deduplicate the request if the first attempt was committed but the response was lost in transit.

// BillingService.java — unsafe mode 1: UUID inside Mono.create() consumer
@Singleton
public class BillingService {

    // Mono.create() is cold: the MonoSink consumer executes on every subscription.
    // retryWhen(Retry.backoff(2, ...)) re-subscribes to the upstream Mono on failure.
    // Each re-subscription re-runs the sink consumer.
    // UUID.randomUUID() inside the consumer generates UUID_B on the first retry.
    // Stripe committed ch_A on attempt 1 before the network error.
    // Retry attempt 2 carries UUID_B — Stripe creates ch_B alongside ch_A.
    public Mono<String> chargeCustomer(String customerId, int amountCents) {
        return Mono.<String>create(sink -> {
            // Inside the sink consumer — executes per Mono.create() subscription.
            // retryWhen() re-subscribes → sink consumer re-runs → UUID.randomUUID() → UUID_B.
            String idempotencyKey = UUID.randomUUID().toString(); // UNSAFE

            ChargeCreateParams params = ChargeCreateParams.builder()
                .setAmount((long) amountCents)
                .setCurrency("usd")
                .setCustomer(customerId)
                .build();

            RequestOptions options = RequestOptions.builder()
                .setIdempotencyKey(idempotencyKey)
                .build();

            // Stripe SDK async: delivers result via CompletableFuture.whenComplete().
            // This is the canonical use case for Mono.create() — wrapping a callback-based async API.
            Charge.createAsync(params, options).whenComplete((charge, error) -> {
                if (error != null) {
                    sink.error(error);
                } else {
                    sink.success(charge.getId());
                }
            });
        })
        .retryWhen(Retry.backoff(2, Duration.ofMillis(200))
            .filter(throwable -> throwable instanceof StripeException stripeEx
                             && isTransient(stripeEx)));
    }

    private boolean isTransient(StripeException ex) {
        return ex instanceof ApiConnectionException || ex instanceof RateLimitException;
    }
}

The call sequence on a transient Stripe network failure:

  1. A caller subscribes to the Mono<String> returned by chargeCustomer().
  2. The Mono.create() sink consumer runs. UUID.randomUUID() generates UUID_A. Charge.createAsync() dispatches the HTTP call asynchronously via the Stripe SDK’s internal thread pool.
  3. Stripe receives the request and begins processing. Stripe commits charge ch_A.
  4. Before the HTTP response completes, the network connection times out or is reset. Charge.createAsync()’s CompletableFuture completes exceptionally with an ApiConnectionException. The sink consumer calls sink.error(error).
  5. The retryWhen() operator catches the error. The filter() predicate returns true for ApiConnectionException. retryWhen() re-subscribes to the upstream Mono.create().
  6. Re-subscription: the Mono.create() sink consumer runs again. UUID.randomUUID() generates UUID_B. Charge.createAsync() dispatches another HTTP call.
  7. Stripe receives the request with UUID_B. It has no record of UUID_B in its idempotency store. Stripe processes it as a new billing intent and commits charge ch_B.
  8. The second attempt succeeds. The caller receives the ch_B charge ID as the result. The customer’s account shows ch_A and ch_B — two charges for the same billing intent. ch_A is unknown to the application.

There is a timing subtlety specific to Mono.create() that does not exist with Mono.fromCallable(): the sink consumer starts an asynchronous operation (Charge.createAsync()) and the result arrives via a callback, not a return value. The retryWhen() re-subscription happens after sink.error() is called, which is after the CompletableFuture completes exceptionally, which is after the Stripe SDK has already made the HTTP call. There is no ambiguity about whether Stripe received the first request — it did. By the time retryWhen() fires, ch_A is committed.

A secondary consideration: the CompletableFuture.whenComplete() callback runs on a thread from the Stripe SDK’s internal ForkJoinPool or ExecutorService, not on a Reactor scheduler. The callback delivers the result to the Reactor pipeline via sink.success() or sink.error(). This works correctly in Reactor — MonoSink is thread-safe — but it means that the completion path crosses a thread boundary between the Stripe SDK thread pool and the Reactor subscriber thread. This is unrelated to the idempotency key problem but is worth noting when debugging timing issues in tests.

The fix is to move UUID.randomUUID() outside the Mono.create() call, to the method scope where it executes once per call to chargeCustomer(). The content-hash variant is preferred because it is also stable across job re-runs:

// BillingService.java — fix for mode 1
@Singleton
public class BillingService {

    public Mono<String> chargeCustomer(String customerId, int amountCents) {
        // Content-hash key: stable across all retryWhen() re-subscriptions AND across job re-runs.
        // chargeCustomer() is called once per billing intent — the key is computed once here.
        // Mono.create() re-subscriptions (via retryWhen) capture this key from the closure.
        final String idempotencyKey = "charge:" + customerId + ":" + amountCents;

        return Mono.<String>create(sink -> {
            // idempotencyKey from closure — same value on every sink consumer re-execution.
            ChargeCreateParams params = ChargeCreateParams.builder()
                .setAmount((long) amountCents)
                .setCurrency("usd")
                .setCustomer(customerId)
                .build();

            RequestOptions options = RequestOptions.builder()
                .setIdempotencyKey(idempotencyKey) // stable key from method scope
                .build();

            Charge.createAsync(params, options).whenComplete((charge, error) -> {
                if (error != null) {
                    sink.error(error);
                } else {
                    sink.success(charge.getId());
                }
            });
        })
        .retryWhen(Retry.backoff(2, Duration.ofMillis(200))
            .filter(throwable -> throwable instanceof StripeException stripeEx
                             && isTransient(stripeEx)));
    }
}

With this fix, attempt 1 sends the content-hash key. If Stripe committed ch_A before the network error, the retry sends the same key. Stripe deduplicates: it recognizes the key and returns ch_A without creating ch_B. One charge, correctly attributed.

A note on the Mono.create() vs Mono.fromCallable() choice: moving to Mono.fromCallable() does not fix the problem unless you also move the UUID outside the callable. It does simplify the structure — Charge.create() (blocking) inside fromCallable() instead of Charge.createAsync().whenComplete() inside create() — but the subscription-execution semantics are the same. Mono.create() is appropriate when you need the non-blocking async path (when the Stripe call shares a thread pool with other async work and you do not want to block a Reactor scheduler thread). In either case, the UUID must be outside the per-subscription factory.

Mode 2: Flux.fromIterable(customerIds).flatMap() batch billing — UUID inside the flatMap() mapper — outer Flux.retryWhen() re-subscribes — all items re-emitted — all mappers re-run — all customers charged twice

The second failure mode involves batch billing: charging a collection of customers in a single reactive pipeline using Flux.fromIterable().flatMap(). This mode is qualitatively more dangerous than mode 1 because the impact is not one duplicate charge on one customer — it is duplicate charges on every customer who was already successfully charged before the failure that triggered Flux.retryWhen(). In a batch of 200 customers, a transient failure on customer 150 can cause customers 1 through 149 to be charged twice.

The failure mechanism is a direct consequence of how Reactor’s Flux.retryWhen() interacts with Flux.fromIterable() and flatMap(). Flux.retryWhen() re-subscribes to the upstream Flux on failure. Re-subscribing to Flux.fromIterable(customerIds) causes it to start from the beginning of the iterable — emitting customer 1 again, then customer 2 again, all the way through. The flatMap() mapper function is called for each emitted item per subscription. Every item re-emitted by the restarted iterable triggers a new flatMap() mapper execution.

UUID inside the flatMap() mapper means that each mapper call generates a new, unique UUID. On the retry re-subscription, customers 1 through 149 each get UUID_B (distinct from UUID_A sent on the first subscription). Stripe has no record of UUID_B for any of these customers — UUID_A was used on the first pass. Stripe processes each UUID_B as a new billing intent and creates ch_B for each customer alongside their already-committed ch_A.

// BatchBillingService.java — unsafe mode 2: UUID inside flatMap() mapper + outer Flux.retryWhen()
@Singleton
public class BatchBillingService {

    private final ReactorHttpClient httpClient;

    // UNSAFE: outer Flux.retryWhen() re-subscribes to the entire Flux on any failure.
    // Re-subscribing to Flux.fromIterable(customerIds) re-emits all customer IDs from index 0.
    // flatMap() mapper executes for each re-emitted item — UUID.randomUUID() inside the mapper
    // generates a new UUID per mapper execution.
    // If customers 1-50 were already charged as ch_A before customer 51 triggered a failure:
    // retryWhen() re-subscribes → customers 1-50 re-emitted → mapper runs for each → UUID_B each
    // → ch_B for each of customers 1-50 alongside their committed ch_A.
    public Flux<BillingResult> chargeBatch(List<String> customerIds, int amountCents,
                                            String billingDate) {
        return Flux.fromIterable(customerIds)
            .flatMap(customerId -> {
                // UUID inside the flatMap() mapper — re-executes per mapper call.
                // Flux.retryWhen() re-subscribes → Flux.fromIterable() re-emits all items →
                // mapper runs for every customer again → UUID.randomUUID() for each → UUID_B.
                String idempotencyKey = UUID.randomUUID().toString(); // UNSAFE

                return Mono.from(httpClient.exchange(
                    HttpRequest.POST("/v1/charges", buildChargeBody(customerId, amountCents))
                        .header("Idempotency-Key", idempotencyKey)
                        .header("Authorization", "Bearer " + stripeApiKey),
                    ChargeResponse.class
                )).map(response -> new BillingResult(customerId, response.body().id()));
            })
            // Outer Flux.retryWhen() retries the entire batch on any failure.
            // "Any failure" includes transient per-customer Stripe errors.
            // This means a transient error on customer N retries all N-1 already-charged customers.
            .retryWhen(Retry.backoff(1, Duration.ofMillis(500))
                .filter(t -> t instanceof StripeTransientException));
    }
}

The call sequence on a transient failure at customer 51 in a 100-customer batch:

  1. Caller subscribes to the Flux returned by chargeBatch().
  2. Flux.fromIterable(customerIds) begins emitting. flatMap() launches concurrent charges for all 100 customers (up to flatMap’s default concurrency of 256).
  3. Customers 1–50 complete successfully: UUID_A for each, ch_A committed per customer.
  4. Customer 51 receives a transient Stripe response: StripeTransientException propagates through the inner Mono into the outer Flux.
  5. The outer Flux.retryWhen() catches the error. The filter predicate matches. retryWhen() re-subscribes to the upstream Flux.
  6. Re-subscription: Flux.fromIterable(customerIds) starts over from index 0. It emits customer 1, customer 2, …, all 100 customers in sequence.
  7. The flatMap() mapper runs for customer 1: UUID.randomUUID() generates UUID_B_1. Stripe sees UUID_B_1, has no record of it, creates ch_B_1 alongside already-committed ch_A_1.
  8. The same happens for customers 2–50. Each gets UUID_B_N and a second charge ch_B_N.
  9. Customer 51 is retried. If the transient error has passed, this attempt succeeds with UUID_B_51 = first Stripe charge for customer 51. Customers 52–100 are also charged for the first time (UUID_B for each).
  10. The application records the results from the retry pass. It has no record of ch_A_1 through ch_A_50. Customers 1–50 each have two charges; customers 52–100 have one each (from the retry pass). Customer 51 has one charge (correctly retried).

There are two distinct problems here that both need to be fixed:

Problem A — UUID regeneration in the mapper: Moving UUID outside the flatMap() mapper eliminates per-mapper-call UUID regeneration. But if retryWhen() still fires on the outer Flux, the entire batch restarts. Even with stable per-customer keys, recharging all 100 customers on a failure of customer 51 is wrong.

Problem B — Batch-level retry instead of per-item retry: The fundamental architectural mistake is applying retry at the batch level (outer Flux.retryWhen()) instead of at the individual charge level (inner Mono.retryWhen() inside the flatMap() mapper). Stripe’s idempotency system can handle per-customer retries safely (with a stable per-customer key) but cannot protect against batch-level retries that re-issue already-committed charges with new keys.

The correct fix has two parts: move UUID to a stable, per-customer scope (content-hash derived from customer ID and billing intent), and move the retry inside the flatMap() mapper to operate per-item rather than per-batch:

// BatchBillingService.java — fix for mode 2: per-item retry with stable content-hash keys
@Singleton
public class BatchBillingService {

    private final ReactorHttpClient httpClient;

    public Flux<BillingResult> chargeBatch(List<String> customerIds, int amountCents,
                                            String billingDate) {
        return Flux.fromIterable(customerIds)
            .flatMap(customerId -> {
                // Content-hash key: stable per customer per billing date.
                // Computed before the inner Mono chain — not inside any deferred operator.
                // If the inner Mono.retryWhen() re-subscribes, the key is the same value.
                // If the batch is re-run tomorrow, the same customer on the same date gets
                // the same key — Stripe deduplicates correctly for billing-job re-runs too.
                final String idempotencyKey = "batch-charge:" + customerId + ":" +
                                              amountCents + ":" + billingDate;

                return Mono.from(httpClient.exchange(
                    HttpRequest.POST("/v1/charges", buildChargeBody(customerId, amountCents))
                        .header("Idempotency-Key", idempotencyKey)
                        .header("Authorization", "Bearer " + stripeApiKey),
                    ChargeResponse.class
                ))
                .map(response -> new BillingResult(customerId, response.body().id()))
                // Retry is per-item, not per-batch:
                // A transient failure on customer 51 retries ONLY customer 51.
                // Customers 1-50 are not re-subscribed or re-charged.
                // The stable idempotencyKey means Stripe deduplicates if ch_A_51
                // was committed before the network error.
                .retryWhen(Retry.backoff(2, Duration.ofMillis(200))
                    .filter(t -> t instanceof StripeTransientException));
            });
        // No outer Flux.retryWhen() — retries happen at the per-customer level above.
    }
}

With per-item retry and a content-hash key, a transient failure on customer 51 triggers a retry of customer 51’s charge only. Customers 1–50 are not re-subscribed. Customer 51’s key is stable across both attempts — if ch_A_51 was committed before the network error, Stripe deduplicates and returns ch_A_51 on the retry. No duplicate charges. No batch restart.

A refinement worth adding for large batches: flatMap(mapper, concurrency) with a controlled concurrency parameter (e.g., 10 or 20) to avoid launching 1,000 concurrent Stripe requests simultaneously, which would hit Stripe’s rate limits and trigger RateLimitException — themselves a source of per-item retries that would consume idempotency key quota. concatMap() (concurrency 1, strictly sequential) is the safest option for billing jobs where idempotency auditability matters most, at the cost of throughput.

One edge case that even the fixed version must handle: if the per-item retryWhen() exhausts all attempts and the inner Mono still fails, that failure propagates up through the outer Flux. Depending on whether you want one failure to abort the entire batch or to continue with remaining customers, you may want to add .onErrorResume(error -> Mono.just(BillingResult.failed(customerId, error))) inside the flatMap() mapper to convert per-customer failures into non-fatal results. The outer batch-level error handling then operates only on the aggregate result, not on any individual Stripe call.

Mode 3: Mono.defer() deferred request construction inside retryWhen(Retry.fixedDelay().filter()) — retryWhen() re-subscribes to the outer Mono including the defer() factory — UUID inside defer() factory regenerates — the custom filter predicate does not prevent re-subscription — UUID_B — ch_B

The third failure mode involves a developer who uses Mono.defer() to build the Micronaut HTTP request and intentionally places UUID.randomUUID() inside the defer() factory. The developer’s reasoning is explicit and internally consistent: “I want each attempt to have its own UUID. Mono.defer() ensures the request is constructed fresh per attempt. A fresh UUID per attempt means each Stripe request is independent.” This reasoning correctly identifies what Mono.defer() does, but it misunderstands what Stripe’s idempotency system requires.

This mode is also distinct from session 290’s onErrorResume() recovery lambda pattern, which involved Mono.defer() inside an error recovery continuation. Here the defer() is the primary upstream publisher, not a recovery path, and the retry is via retryWhen(Retry.fixedDelay().filter()) with a custom exception-type predicate.

Stripe’s documentation is explicit: the idempotency key identifies a request for an operation, not a single HTTP attempt. If you send the same operation (charge $99 to customer_X) twice with the same key, Stripe returns the result of the first operation — it will not charge the customer again. This is precisely the guarantee you want during a network-level retry: Stripe’s commitment is that if you made the exact same API call earlier, it is safe to make it again with the same key, because Stripe will return the cached result without re-executing the operation. Using a new UUID on each retry disables this guarantee. Each retry carries a key Stripe has never seen, and Stripe treats it as a new, independent billing intent.

// PaymentClient.java — unsafe mode 3: UUID inside Mono.defer() factory
@Singleton
public class PaymentClient {

    private final ReactorHttpClient httpClient;

    // Mono.defer() is cold: the Supplier<Mono<T>> executes on every subscription.
    // retryWhen(Retry.fixedDelay(2, ...).filter(predicate)) re-subscribes to the upstream
    // Mono on failures where predicate returns true.
    // Re-subscribing to Mono.defer() re-executes the factory supplier.
    // UUID.randomUUID() inside the supplier generates UUID_B on the first retry.
    // The filter() predicate determines WHICH failures trigger retry — it does NOT prevent
    // the defer() factory from re-evaluating when retry does fire.
    public Mono<String> createCharge(String customerId, long amountCents, String currency) {
        return Mono.defer(() -> {
            // Inside Mono.defer() factory — executes per subscription.
            // Developer intention: "fresh request per attempt, including a fresh UUID."
            // Stripe's intention: "same UUID for the same billing intent across all attempts."
            String idempotencyKey = UUID.randomUUID().toString(); // UNSAFE

            ChargeRequestBody body = new ChargeRequestBody(customerId, amountCents, currency);

            return Mono.from(httpClient.exchange(
                HttpRequest.POST("/v1/charges", body)
                    .header("Authorization", "Bearer " + stripeApiKey)
                    .header("Idempotency-Key", idempotencyKey),
                ChargeResponse.class
            )).map(response -> response.body().id());
        })
        .retryWhen(
            Retry.fixedDelay(2, Duration.ofMillis(300))
                // Custom filter: only retry on connection errors and 5xx responses.
                // 4xx errors (invalid card, insufficient funds) are not retried.
                // This filter limits WHICH exceptions trigger retryWhen().
                // It does NOT prevent defer() from re-evaluating UUID when retry fires.
                .filter(throwable -> {
                    if (throwable instanceof ConnectException) return true;
                    if (throwable instanceof HttpClientResponseException hcre) {
                        return hcre.getStatus().getCode() >= 500;
                    }
                    return false;
                })
        );
    }
}

The developer’s mental model of Mono.defer() is correct on its own terms: the factory supplier does execute per subscription, and each subscription does start a fresh HttpRequest with a fresh UUID. But the developer’s mental model of retry is wrong: they are treating each Reactor retry subscription as an independent billing attempt that should have its own UUID, when Stripe’s model treats all HTTP attempts for the same billing operation as retries of the same billing intent that must share a UUID.

The filter() predicate in Retry.fixedDelay().filter() adds a superficially plausible mental model: “I only retry on connection errors and 5xx errors, not on 4xx errors. So I know the first attempt failed before Stripe processed it — the payment was not taken.” This reasoning is subtly wrong for connection errors: a ConnectException means the TCP connection failed, but ApiConnectionException from the Stripe SDK means the HTTP request was sent and the response was not received — the latter does not mean Stripe did not process the charge. Even restricting to genuine pre-connection failures (ConnectException before any bytes were sent), the code also retries on HttpClientResponseException with 5xx status codes, which means Stripe did receive and process the request but returned an error response — and the underlying charge may or may not have been committed depending on which 5xx code was returned. For 500 errors, Stripe recommends using the same idempotency key on the retry to avoid duplicate charges. Sending UUID_B on a 500 retry creates exactly the duplicate this predicate was meant to prevent.

The fix has two parts. First, move UUID generation outside Mono.defer() to the method scope — use a content-hash key. Second, keep Mono.defer() if you want a fresh HttpRequest object per subscription (this is a legitimate and useful pattern for other reasons, such as ensuring request immutability and avoiding header state leakage across retries), but capture the stable key from the enclosing scope:

// PaymentClient.java — fix for mode 3
@Singleton
public class PaymentClient {

    private final ReactorHttpClient httpClient;

    public Mono<String> createCharge(String customerId, long amountCents, String currency) {
        // Content-hash key: stable across all retryWhen() re-subscriptions.
        // Computed at method scope — not inside Mono.defer() or any cold factory.
        // Same customerId + amountCents + currency always produces the same key.
        // Stripe deduplicates correctly on retry: if ch_A was committed before
        // the network error, the retry returns ch_A without creating ch_B.
        final String idempotencyKey = "charge:" + customerId + ":" + amountCents + ":" + currency;

        return Mono.defer(() -> {
            // Mono.defer() is still useful here: it ensures HttpRequest is built
            // fresh per subscription, not shared across concurrent subscribers.
            // But the idempotencyKey is captured from method scope — same value per subscription.
            ChargeRequestBody body = new ChargeRequestBody(customerId, amountCents, currency);

            return Mono.from(httpClient.exchange(
                HttpRequest.POST("/v1/charges", body)
                    .header("Authorization", "Bearer " + stripeApiKey)
                    .header("Idempotency-Key", idempotencyKey), // stable key from closure
                ChargeResponse.class
            )).map(response -> response.body().id());
        })
        .retryWhen(
            Retry.fixedDelay(2, Duration.ofMillis(300))
                .filter(throwable -> {
                    if (throwable instanceof ConnectException) return true;
                    if (throwable instanceof HttpClientResponseException hcre) {
                        return hcre.getStatus().getCode() >= 500;
                    }
                    return false;
                })
        );
    }
}

With this fix, all Mono.defer() subscriptions (on initial subscription and on all retryWhen() re-subscriptions) use the same idempotencyKey value captured from the method scope. The HttpRequest object is built fresh per subscription (no sharing), but the idempotency key it carries is stable. Stripe deduplicates correctly.

The interaction between the filter() predicate and idempotency deserves a final note. The predicate gates which exceptions trigger retry — ConnectException and 5xx responses. For ConnectException (TCP connection refused before any bytes sent), there is a strong argument that UUID_A was never received by Stripe and a new UUID would be fine. But there are two reasons to still use a stable key here: first, Stripe’s documentation recommends a stable key even when you are confident the first attempt was not processed, because the key provides an audit trail; second, the content-hash key costs nothing and removes the need to reason about which exception types are “pre-receipt” versus “post-receipt.” A stable content-hash key is the right default. If you are retrying a 5xx response (Stripe received and processed the request), a stable key is not merely a recommendation — it is the only correct behavior.

Cross-mode comparison: cold-publisher type, retry operator, UUID boundary, and blast radius

Mode Cold-publisher factory Retry operator Scope where UUID lives (unsafe) What re-executes per retry Blast radius
1: Mono.create() sink consumer Mono.create(MonoSink consumer) Mono.retryWhen(Retry.backoff) Inside the MonoSink consumer lambda The entire sink consumer body (UUID, request build, async dispatch) One charge: ch_B for the single customer
2: Flux.flatMap() batch Flux.fromIterable() + flatMap(mapper) Outer Flux.retryWhen(Retry.backoff) Inside the flatMap() mapper function The entire Flux from the beginning: all items re-emitted, all mappers re-run All already-charged customers in the batch: ch_B for every customer 1 through N-1
3: Mono.defer() + filtered retry Mono.defer(Supplier factory) Mono.retryWhen(Retry.fixedDelay().filter()) Inside the Mono.defer() factory supplier The defer() factory (UUID, request build, HTTP call) One charge: ch_B for the single customer (for each filtered-in retry)

The batch mode (mode 2) is the most dangerous because a single transient failure on any item in the batch triggers a full restart that re-charges all already-billed customers. The single-charge modes (1 and 3) produce at most two charges per customer per billing cycle. But in all three cases the root cause is identical: UUID generation is inside a cold-publisher factory that re-executes per subscription, and the retry operator causes re-subscription.

Where the stable key must live in each mode

The universal fix is to compute the idempotency key at a scope that executes exactly once per billing intent, regardless of how many Reactor subscriptions the downstream pipeline creates:

In all three cases, the content-hash key — a deterministic string derived from the billing intent parameters — is preferred over UUID at the correct scope. A UUID at the correct scope (method scope) prevents retry-induced key regeneration, but it does not prevent re-key on billing job re-runs: if the batch job fails halfway through and is restarted, a UUID computed at the start of the job run is different from the UUID in the original run. A content-hash key derived from (customerId, amountCents, billingDate) is the same in both runs, so Stripe deduplicates correctly even across job restarts.

Testing each mode with @MicronautTest, WireMock, and StepVerifier

Testing that idempotency keys are stable across retry requires capturing all Idempotency-Key header values sent to Stripe across all retry attempts and asserting that they are the same value. WireMock’s request journaling makes this straightforward.

// BillingServiceTest.java — testing mode 1 (Mono.create() + retryWhen)
@MicronautTest
class BillingServiceTest {

    @Inject
    BillingService billingService;

    private WireMockServer wireMock;

    @BeforeEach
    void setUp() {
        wireMock = new WireMockServer(wireMockConfig().port(18443));
        wireMock.start();
        // Configure Micronaut test to route Stripe calls to WireMock
    }

    @AfterEach
    void tearDown() { wireMock.stop(); }

    @Test
    void chargeCustomer_sendsStableIdempotencyKeyAcrossRetries() {
        // WireMock scenario: first call returns 500, second call returns 200.
        wireMock.stubFor(post(urlEqualTo("/v1/charges"))
            .inScenario("stripe-retry")
            .whenScenarioStateIs(STARTED)
            .willReturn(serverError())
            .willSetStateTo("after-first"));

        wireMock.stubFor(post(urlEqualTo("/v1/charges"))
            .inScenario("stripe-retry")
            .whenScenarioStateIs("after-first")
            .willReturn(okJson("{\"id\":\"ch_abc123\",\"object\":\"charge\"}")));

        StepVerifier.create(billingService.chargeCustomer("cus_test", 9900))
            .expectNext("ch_abc123")
            .verifyComplete();

        // Capture all requests sent to WireMock — includes both attempts.
        List<LoggedRequest> requests = wireMock.findAll(
            postRequestedFor(urlEqualTo("/v1/charges")));

        assertThat(requests).hasSize(2);

        // Extract idempotency keys from both requests.
        String keyAttempt1 = requests.get(0).getHeader("Idempotency-Key");
        String keyAttempt2 = requests.get(1).getHeader("Idempotency-Key");

        // Assert: both attempts used the same idempotency key.
        // If UUID was inside Mono.create(), these would be different (UUID_A, UUID_B).
        // The fix (content-hash at method scope) makes them identical.
        assertThat(keyAttempt1).isEqualTo(keyAttempt2);
        // Also assert the key is deterministic from the input parameters.
        assertThat(keyAttempt1).isEqualTo("charge:cus_test:9900");
    }
}
// BatchBillingServiceTest.java — testing mode 2 (Flux.flatMap() batch + per-item retry)
@MicronautTest
class BatchBillingServiceTest {

    @Inject
    BatchBillingService batchBillingService;

    @Test
    void chargeBatch_perItemRetry_doesNotRechargePreviousCustomers() {
        List<String> customerIds = List.of("cus_001", "cus_002", "cus_003");

        // WireMock: cus_001 succeeds immediately, cus_002 fails once then succeeds,
        // cus_003 succeeds immediately.
        wireMock.stubFor(post(urlEqualTo("/v1/charges"))
            .withHeader("Authorization", containing("cus_001-key"))
            .willReturn(okJson("{\"id\":\"ch_001\"}")));

        wireMock.stubFor(post(urlEqualTo("/v1/charges"))
            .withHeader("Idempotency-Key", equalTo("batch-charge:cus_002:9900:2026-10-03"))
            .inScenario("cus_002-retry").whenScenarioStateIs(STARTED)
            .willReturn(serverError()).willSetStateTo("after-first"));

        wireMock.stubFor(post(urlEqualTo("/v1/charges"))
            .withHeader("Idempotency-Key", equalTo("batch-charge:cus_002:9900:2026-10-03"))
            .inScenario("cus_002-retry").whenScenarioStateIs("after-first")
            .willReturn(okJson("{\"id\":\"ch_002\"}")));

        wireMock.stubFor(post(urlEqualTo("/v1/charges"))
            .withHeader("Idempotency-Key", equalTo("batch-charge:cus_003:9900:2026-10-03"))
            .willReturn(okJson("{\"id\":\"ch_003\"}")));

        StepVerifier.create(batchBillingService.chargeBatch(customerIds, 9900, "2026-10-03"))
            .expectNextCount(3)
            .verifyComplete();

        // Assert: cus_001 and cus_003 each received exactly ONE Stripe request.
        // (Per-item retry on cus_002 must not re-trigger cus_001 or cus_003.)
        List<LoggedRequest> allRequests = wireMock.findAll(
            postRequestedFor(urlEqualTo("/v1/charges")));

        long cus001Requests = allRequests.stream()
            .filter(r -> r.getHeader("Idempotency-Key")
                          .equals("batch-charge:cus_001:9900:2026-10-03"))
            .count();
        long cus002Requests = allRequests.stream()
            .filter(r -> r.getHeader("Idempotency-Key")
                          .equals("batch-charge:cus_002:9900:2026-10-03"))
            .count();

        // cus_001: exactly 1 request (no retry — only cus_002 retried).
        assertThat(cus001Requests).isEqualTo(1);
        // cus_002: exactly 2 requests (1 failure + 1 retry, both with the same key).
        assertThat(cus002Requests).isEqualTo(2);

        // All keys stable: set of unique keys for cus_002 == 1 (same key both times).
        Set<String> cus002Keys = allRequests.stream()
            .filter(r -> r.getHeader("Idempotency-Key")
                          .equals("batch-charge:cus_002:9900:2026-10-03"))
            .map(r -> r.getHeader("Idempotency-Key"))
            .collect(Collectors.toSet());
        assertThat(cus002Keys).hasSize(1);
    }
}

The test for mode 2 focuses on the per-item vs batch-level retry distinction: it asserts that cus_001 received exactly one Stripe request even though cus_002 required a retry. This directly validates that per-item retry does not restart the entire batch. The key stability assertion (same key on both cus_002 attempts) validates that the content-hash key is stable across per-item retries.

For mode 3 (filtered retryWhen), the test structure is similar to mode 1 but should also include a negative test: verify that a 4xx response (filtered out by the predicate) is NOT retried and the error propagates without a second Stripe request. This tests that the filter predicate is working correctly and that the application does not silently swallow non-retryable errors.

The developer’s mental model gap in all three modes

Each of the three modes has a developer mental model that is internally consistent but conflicts with Stripe’s idempotency contract:

The common thread: the correct mental model is that the idempotency key belongs to the billing intent, not to the HTTP attempt or the Reactor subscription. A billing intent is defined at the application level (charge customer X amount Y for order Z) and persists across all retry attempts, all Reactor subscriptions, and all job re-runs that implement the same intent. The key should be computed from the business parameters that define the billing intent, before any code that might execute more than once per intent invocation.

The same structural analysis applies to other reactive Micronaut HTTP client patterns: Micronaut Data reactive @Transactional methods with retryWhen(), and Quarkus reactive @Transactional with Panache. The reactive framework differs, but the fix is always the same: compute the idempotency key at the billing-intent scope, outside every per-subscription and per-method-body boundary.

Keybrake catches duplicate Stripe keys before they charge your customers twice

Keybrake is a scoped API-key proxy that sits between your application and Stripe. It enforces idempotency key uniqueness per billing intent — rejecting retries that arrive with a new key for an already-committed charge, and logging every key collision for audit. Join the waitlist.