Spring Data MongoDB and Stripe Integration: How @Retryable + @Transactional(ReactiveMongoTransactionManager) Re-executes UUID Generation on Each proceed() Call, ReactiveMongoTemplate.inTransaction() Re-invokes the Callback on retryWhen() Re-subscription and Re-evaluates Mono.defer{} UUID, and ReactiveMongoRepository.findAll() Cold Flux Re-queries MongoDB on Outer retryWhen() and Recharges the Entire Batch
Spring Data MongoDB’s reactive stack introduces three idempotency failure patterns that are structurally similar to the JPA equivalents but differ in how ReactiveMongoTransactionManager manages reactive sessions, how ReactiveMongoTemplate.inTransaction() re-invokes its callback on retryWhen() re-subscription, and how ReactiveMongoRepository’s derived query methods produce cold Flux publishers backed by fresh MongoDB cursors on every new subscription. In each case, a UUID computed inside the retried code boundary produces UUID_B on the second attempt, and Stripe treats UUID_B as a distinct payment intent — ch_B alongside committed ch_A.
This post is structurally distinct from the Spring Data JPA post (which covers EntityManager lifecycle under @Retryable + @Transactional, JPA optimistic locking retry loops, and @TransactionalEventListener + @Retryable) and from the Spring WebFlux post (which covers the basic Mono.defer{} + retryWhen() pattern without a data layer). The failure modes here are specific to how ReactiveMongoTransactionManager participates in Spring’s AOP proxy chain, how ReactiveMongoTemplate.inTransaction() wraps a callback function and what happens to that function under reactive retry, and how MongoDB’s reactive repository implementation produces cold publishers for every derived query method.
Background: Spring Data MongoDB reactive stack, ReactiveMongoTransactionManager, and cold Flux publishers
Spring Data MongoDB’s reactive support is built on top of the MongoDB Reactive Streams Java Driver. The driver implements the Reactive Streams specification; Spring Data wraps it with Project Reactor types (Mono<T> and Flux<T>). Every ReactiveMongoRepository query method — including findAll(), findAllByActiveTrue(), findById(), and any derived query method — returns a cold publisher. A cold publisher is one that performs no work and holds no connection until a subscriber subscribes. Each new subscription triggers a fresh MongoDB query: a new network round-trip, a new cursor on the MongoDB server, and a new result stream from the beginning.
This is the correct default behavior for a reactive data layer. It means that if you do repository.findAll().subscribe(consumer1) and then later repository.findAll().subscribe(consumer2), each subscriber gets a fresh result set reflecting the state of the collection at subscription time. No shared mutable cursor, no position state shared between subscribers.
The problem arises when a retry operator is placed on the outer cold Flux. Reactive retry operators work by re-subscribing the upstream publisher on failure. When the upstream is a cold MongoDB repository query, re-subscription re-queries the database. The entire document stream starts over. Any processing applied inside a flatMap on that stream — including UUID generation and Stripe charges — executes again for every document that is re-emitted.
Spring Data MongoDB supports multi-document ACID transactions starting with MongoDB 4.0 and requiring a replica set (or sharded cluster). Transactions are managed via ReactiveMongoTransactionManager, which implements Spring’s ReactiveTransactionManager interface. When you annotate a method with @Transactional and configure the transaction manager bean as ReactiveMongoTransactionManager, Spring’s transaction interceptor opens a ClientSession on the reactive MongoDB driver, associates it with the current reactive context (Context in Project Reactor), and ensures that all repository operations executed within that method use the same session. The transaction commits when the reactive pipeline completes successfully and rolls back on failure.
The interaction with @Retryable follows the same AOP proxy ordering as the JPA case. Spring Retry’s RetryOperationsInterceptor is configured with order Integer.MAX_VALUE - 5 (value: 2147483642) by default. Spring’s TransactionInterceptor is configured with order Integer.MAX_VALUE (value: 2147483647). A lower order number means the interceptor runs earlier in the chain — it is the outer wrapper. So @Retryable is always the outer interceptor when both annotations are present on the same method. On retry, Spring Retry calls MethodInvocation.proceed(). This re-invokes the interceptor chain starting from the @Transactional interceptor inward. A new reactive MongoDB ClientSession is opened. The method body re-executes from its first line. Any UUID.randomUUID() call at method entry generates UUID_B.
Stripe’s idempotency contract: a POST to any Stripe mutating endpoint with an Idempotency-Key header deduplicates requests carrying the same key per API key for 24 hours. Two distinct keys for the same customer, amount, and billing period are two distinct payment intents. Both are charged.
Mode 1: @Retryable + @Transactional(transactionManager = “reactiveMongoTransactionManager”) — proceed() re-executes the full method body including UUID generation
The developer is building a billing service on Spring Data MongoDB. They want multi-document ACID guarantees when writing billing records to MongoDB (their primary store is MongoDB, not a relational DB), so they configure ReactiveMongoTransactionManager and annotate their billing method with @Transactional. They also want retry semantics for transient Stripe network errors, so they add @Retryable.
The method computes the Idempotency-Key for Stripe at the top of the method body before any reactive pipeline is assembled. The developer’s mental model: “UUID is the first thing I do — it’s a plain Java expression, not inside any reactive pipeline or transaction lambda. @Retryable retries the service-level operation, which means it retries the Stripe HTTP call. @Transactional with ReactiveMongoTransactionManager manages the MongoDB reactive session; when the transaction rolls back on a Stripe failure, it undoes the MongoDB write but the UUID was already computed before any transaction started. Retrying means re-entering the transaction, not re-generating the UUID that was set before the transaction.”
This mental model is wrong about what @Retryable retries. Spring Retry’s interceptor wraps the full method invocation at the AOP proxy level. When it retries, it calls MethodInvocation.proceed(), which invokes the entire interceptor chain again — starting with the @Transactional interceptor which opens a new reactive MongoDB session, and then re-executing the method body from line 1. The UUID computation on line 1 is not outside the retry scope. It is the first line of the method body that proceed() re-invokes.
// BillingService.java — unsafe Mode 1: @Retryable outer AOP + @Transactional inner
@Service
public class MongoReactiveBillingService {
private final CustomerRepository customerRepository;
private final BillingAttemptRepository billingAttemptRepository;
private final StripeWebClient stripeClient;
// Developer's reasoning:
// "UUID.randomUUID() is on line 1, outside the reactive pipeline and outside the
// @Transactional boundary. @Retryable retries the Stripe HTTP call. @Transactional
// with ReactiveMongoTransactionManager manages the MongoDB reactive session lifecycle.
// UUID is generated before any transaction starts — it's not inside the retry scope."
//
// The problem:
// AOP proxy order: @Retryable (order MAX-5) is outer; @Transactional (order MAX) is inner.
// On a Stripe transient failure, Spring Retry calls MethodInvocation.proceed().
// proceed() re-invokes: @Transactional interceptor opens new ClientSession → method body runs.
// Line 1 of the method body: UUID.randomUUID() → UUID_B. ← NEW UUID
//
// Timeline:
// Attempt 1:
// @Retryable intercepts → calls proceed() → @Transactional opens ClientSession1/Tx1.
// UUID_A = UUID.randomUUID().toString() ← generated here
// billingAttemptRepository.save(BillingAttempt{key: UUID_A}) → persists in Tx1.
// stripeClient.charge(UUID_A) → Stripe processes ch_A → returns 500 (network glitch).
// Stripe side effect: ch_A IS committed (Stripe 500 on PaymentIntent is ambiguous).
// @Transactional rolls back Tx1 (MongoDB write rolled back).
// @Retryable catches the exception, increments retry count.
//
// Attempt 2:
// @Retryable calls proceed() again → @Transactional opens ClientSession2/Tx2.
// UUID_B = UUID.randomUUID().toString() ← NEW UUID on line 1
// billingAttemptRepository.save(BillingAttempt{key: UUID_B}) → persists in Tx2.
// stripeClient.charge(UUID_B) → Stripe processes ch_B. ← DUPLICATE
//
// Result: ch_A and ch_B both committed. Customer charged twice.
@Retryable(
retryFor = { StripeNetworkException.class, StripeTimeoutException.class },
maxAttempts = 3,
backoff = @Backoff(delay = 200, multiplier = 2)
)
@Transactional("reactiveMongoTransactionManager")
public Mono<String> chargeCustomer(String customerId, String stripeCustomerId,
long amountCents, String billingPeriod) {
// Developer places UUID generation here, reasoning it is "outside the retry loop."
// In reality, this is the first statement of the method body that proceed() re-invokes.
String idempotencyKey = UUID.randomUUID().toString();
return billingAttemptRepository
.save(new BillingAttempt(customerId, amountCents, billingPeriod, idempotencyKey))
.flatMap(attempt ->
stripeClient.createPaymentIntent(stripeCustomerId, amountCents,
idempotencyKey))
.flatMap(paymentIntentId ->
billingAttemptRepository.markSucceeded(customerId, billingPeriod,
paymentIntentId));
}
}
There is a MongoDB-specific nuance that makes this failure mode more dangerous than it looks. Reactive MongoDB transactions require a ClientSession that is propagated through the Project Reactor context. When @Transactional rolls back Tx1 due to the Stripe exception, it closes ClientSession1. On retry, @Transactional opens ClientSession2 and puts it into a new reactive context. The MongoDB write of BillingAttempt{key: UUID_A} was rolled back (never committed to the replica set). The Stripe charge ch_A, however, is committed — Stripe is an external HTTP call and is not part of the MongoDB transaction. When attempt 2 proceeds with UUID_B, ch_B is added. The developer may look at their MongoDB audit collection after the incident and see only BillingAttempt{key: UUID_B, status: succeeded} — the UUID_A record was rolled back and is not visible. They have no MongoDB evidence that the original charge with UUID_A succeeded, which makes debugging harder. The Stripe dashboard shows both ch_A and ch_B; the MongoDB audit log shows only ch_B. The audit gap between the two systems is the detection signal.
The fix: content-hash UUID in a non-@Retryable facade
The fix is identical in structure to the JPA fix: separate UUID computation from the retried method invocation. A facade method that does not carry @Retryable computes a deterministic content-hash idempotency key from the billing inputs. It passes the key as a stable parameter to the inner method that carries @Retryable and @Transactional. The inner method receives the same key regardless of how many times proceed() re-invokes it.
// BillingService.java — safe Mode 1 fix: content-hash UUID in facade, stable param
@Service
public class MongoReactiveBillingService {
private final BillingAttemptRepository billingAttemptRepository;
private final StripeWebClient stripeClient;
// Public API: facade computes stable key, no @Retryable.
public Mono<String> chargeCustomer(String customerId, String stripeCustomerId,
long amountCents, String billingPeriod) {
// Content-hash key: deterministic from billing inputs.
// Same customer + billingPeriod + amount always produces the same key within
// a 24-hour Stripe idempotency window — safe for any number of retries.
String idempotencyKey = Hashing.sha256()
.hashString(customerId + "|" + billingPeriod + "|" + amountCents,
StandardCharsets.UTF_8)
.toString()
.substring(0, 36);
return doChargeWithKey(customerId, stripeCustomerId, amountCents,
billingPeriod, idempotencyKey);
}
// Inner method: receives stable key as parameter.
// @Retryable + @Transactional are safe here because the key comes from the caller.
@Retryable(
retryFor = { StripeNetworkException.class, StripeTimeoutException.class },
maxAttempts = 3,
backoff = @Backoff(delay = 200, multiplier = 2)
)
@Transactional("reactiveMongoTransactionManager")
Mono<String> doChargeWithKey(String customerId, String stripeCustomerId,
long amountCents, String billingPeriod,
String idempotencyKey) {
// idempotencyKey is a method parameter — proceed() passes the same value on retry.
return billingAttemptRepository
.save(new BillingAttempt(customerId, amountCents, billingPeriod, idempotencyKey))
.flatMap(attempt ->
stripeClient.createPaymentIntent(stripeCustomerId, amountCents,
idempotencyKey))
.flatMap(paymentIntentId ->
billingAttemptRepository.markSucceeded(customerId, billingPeriod,
paymentIntentId));
}
}
The MongoDB-specific consideration here is that doChargeWithKey must be called through the Spring AOP proxy to get the @Retryable and @Transactional interceptors. If chargeCustomer calls this.doChargeWithKey() directly, Spring’s AOP proxy is bypassed (same limitation as in the JPA case). The standard Spring idiom is to self-inject the bean: declare @Autowired private MongoReactiveBillingService self in the same class and call self.doChargeWithKey(). Alternatively, extract doChargeWithKey into a separate @Service class that carries the @Retryable annotation, and inject that class into the facade.
There is also the question of what happens to the MongoDB BillingAttempt write on repeated retries with the same content-hash key. On attempt 1, the attempt record is saved with the content-hash key. If attempt 1 rolls back the transaction (because Stripe returned a transient error and the MongoDB write was part of the same transaction), then on attempt 2, the save() is fresh — the previous record does not exist. If you want an idempotent MongoDB write as well (to handle the case where the MongoDB write on attempt 1 succeeded but Stripe failed), you should use upsert with the content-hash key as the unique field: find-or-create by idempotency key before proceeding to the Stripe call. This ensures that if attempt 1 fully committed (MongoDB write succeeded, Stripe succeeded) and the caller retries due to a network timeout on the response, the save() on attempt 2 is a no-op.
Mode 2: ReactiveMongoTemplate.inTransaction() with Mono.defer{} UUID inside the callback — retryWhen() re-invokes the callback Function
The developer wants fine-grained control over MongoDB transaction scope and switches from the declarative @Transactional annotation to the programmatic ReactiveMongoTemplate.inTransaction() API. This API takes a Function<ReactiveMongoOperations, Publisher<T>> (a callback lambda that receives a ReactiveMongoOperations bound to the transaction session) and returns a Mono<T> that, when subscribed, opens a transaction, invokes the callback to get a Publisher<T>, subscribes to it, and commits or rolls back based on the result.
The developer wants to be reactive-idiomatic with UUID generation. They know that eager UUID computation in a reactive pipeline can cause issues (they’ve read about Mono.just(UUID.randomUUID()) being evaluated at assembly time, which is not a problem for UUID but can be a problem for other computations). They use Mono.defer{} inside the callback lambda to defer UUID computation to subscription time, thinking this makes the UUID generation “reactive-correct” — it will run at the right point in the reactive lifecycle, i.e., when the transaction is open and active.
They then apply retryWhen() to the outer Mono returned by inTransaction(), intending to retry transient Stripe failures. The problem: retryWhen() re-subscribes the upstream publisher on failure. The upstream publisher is the Mono returned by inTransaction(). When inTransaction()’s Mono is re-subscribed, the inTransaction() operator must produce a new transaction pipeline. It calls the callback Function again to get a new Publisher<T>. The callback lambda body re-executes. Any Mono.defer{} inside the lambda is re-evaluated on the new subscription. UUID_B is generated.
// BillingService.java — unsafe Mode 2: inTransaction() + Mono.defer{} UUID + retryWhen()
@Service
public class MongoProgrammaticBillingService {
private final ReactiveMongoTemplate mongoTemplate;
private final StripeWebClient stripeClient;
// Developer's reasoning:
// "ReactiveMongoTemplate.inTransaction() takes a callback Function and runs it inside a
// MongoDB multi-document transaction. I use Mono.defer{} to defer UUID generation to
// subscription time — reactive-correct, runs when the transaction is open.
// retryWhen() retries the Stripe HTTP call inside the pipeline. The inTransaction()
// block has already committed by the time retryWhen() fires on a Stripe failure —
// the callback Function is done; retryWhen() only re-runs the Stripe-specific Mono."
//
// The problem:
// retryWhen() is applied to the Mono returned by inTransaction().
// On a Stripe failure, retryWhen() re-subscribes the entire Mono returned by inTransaction().
// inTransaction()'s Mono re-invokes the callback Function to get a new Publisher.
// The callback lambda body re-executes.
// Mono.defer{ UUID.randomUUID() } inside the lambda re-evaluates its factory.
// UUID_B is generated. Stripe receives UUID_B. ch_B alongside committed ch_A.
//
// Timeline:
// Subscription 1:
// retryWhen() subscribes to inTransaction()'s Mono.
// inTransaction() calls the callback Function → callback body runs.
// Mono.defer{ ... } factory executes: UUID_A = UUID.randomUUID().toString()
// BillingAttempt saved with UUID_A inside the transaction.
// stripeClient.charge(UUID_A) → Stripe returns 500. ch_A committed by Stripe.
// inTransaction()'s Mono emits failure.
// retryWhen() catches the failure, decides to retry.
//
// Subscription 2:
// retryWhen() re-subscribes to inTransaction()'s Mono.
// inTransaction() calls the callback Function again → callback body runs again.
// Mono.defer{ ... } factory executes again: UUID_B = UUID.randomUUID().toString() ← NEW
// BillingAttempt saved with UUID_B inside the new transaction.
// stripeClient.charge(UUID_B) → ch_B. ← DUPLICATE
public Mono<String> chargeCustomer(String customerId, String stripeCustomerId,
long amountCents, String billingPeriod) {
return mongoTemplate.inTransaction().execute(ops ->
// Mono.defer{} defers UUID evaluation to subscription time.
// Developer's intention: "UUID runs inside the open transaction context."
// The problem: this is inside the callback Function — it re-runs when
// inTransaction()'s Mono is re-subscribed by retryWhen().
Mono.defer(() -> {
String idempotencyKey = UUID.randomUUID().toString();
return ops.save(new BillingAttempt(customerId, amountCents,
billingPeriod, idempotencyKey))
.flatMap(attempt ->
stripeClient.createPaymentIntent(stripeCustomerId,
amountCents, idempotencyKey));
})
).retryWhen(Retry.backoff(3, Duration.ofMillis(200))
.filter(e -> e instanceof StripeNetworkException));
}
}
The developer misconception has a specific form that is distinct from the JPA or Quarkus CDI cases. In those cases, the developer correctly identifies that @Transactional and @Retryable are separate AOP interceptors and reasons about which one is outer. In this case, the developer is reasoning about the reactive operator level: they believe inTransaction() executes the callback once and commits a discrete transaction, and that retryWhen() only retries the portion of the pipeline after the transaction committed. This reflects a misunderstanding of how cold publishers and retry operators interact. retryWhen() does not retry “from the Stripe call.” It re-subscribes the entire upstream publisher — which is the Mono produced by inTransaction(), which re-invokes the entire callback.
The fix: UUID computed outside the retryWhen() scope
The fix is to compute the UUID before the retryWhen() retry boundary. The simplest form: compute UUID as an eager local variable in the method body, outside the inTransaction() callback, and capture it as a final variable in the lambda.
// BillingService.java — safe Mode 2 fix: UUID outside retryWhen() scope
@Service
public class MongoProgrammaticBillingService {
private final ReactiveMongoTemplate mongoTemplate;
private final StripeWebClient stripeClient;
public Mono<String> chargeCustomer(String customerId, String stripeCustomerId,
long amountCents, String billingPeriod) {
// UUID computed at method-body assembly time — outside inTransaction() and retryWhen().
// Content-hash: deterministic for the same billing inputs within 24-hour Stripe window.
final String idempotencyKey = Hashing.sha256()
.hashString(customerId + "|" + billingPeriod + "|" + amountCents,
StandardCharsets.UTF_8)
.toString()
.substring(0, 36);
// idempotencyKey is a final local variable — captured by the lambda.
// When retryWhen() re-subscribes inTransaction()'s Mono, the callback Function
// is re-invoked, but idempotencyKey is not recomputed — it's a captured final.
return mongoTemplate.inTransaction().execute(ops ->
ops.save(new BillingAttempt(customerId, amountCents, billingPeriod,
idempotencyKey))
.flatMap(attempt ->
stripeClient.createPaymentIntent(stripeCustomerId,
amountCents, idempotencyKey))
).retryWhen(Retry.backoff(3, Duration.ofMillis(200))
.filter(e -> e instanceof StripeNetworkException));
}
}
The captured idempotencyKey is a final local variable in the enclosing method. When the callback Function lambda is re-invoked by inTransaction() on retryWhen() re-subscription, the lambda body re-executes — but the variable reference resolves to the same captured value. No new UUID is generated on retry.
There is an additional consideration specific to inTransaction() in a retry scenario. On attempt 1, the MongoDB transaction for the BillingAttempt save may have committed before Stripe returned the 500 (if the save and Stripe call are sequential and the save succeeded). On attempt 2, inTransaction() opens a new transaction and calls save() again with the same idempotency key. If your BillingAttempt collection has a unique index on the idempotency key field, the second save() will throw a DuplicateKeyException. This is actually desirable behavior for an idempotent write — you want to detect that the operation already ran. Handle it explicitly: catch DuplicateKeyException from the save() step and treat it as a successful idempotent repeat. The Stripe charge, having been made with the same idempotency key, will be deduplicated by Stripe automatically if it previously succeeded.
// BillingService.java — safe Mode 2 with idempotent MongoDB write handling
public Mono<String> chargeCustomer(String customerId, String stripeCustomerId,
long amountCents, String billingPeriod) {
final String idempotencyKey = Hashing.sha256()
.hashString(customerId + "|" + billingPeriod + "|" + amountCents,
StandardCharsets.UTF_8)
.toString()
.substring(0, 36);
return mongoTemplate.inTransaction().execute(ops ->
ops.upsert(
Query.query(Criteria.where("idempotencyKey").is(idempotencyKey)),
Update.update("customerId", customerId)
.set("amountCents", amountCents)
.set("billingPeriod", billingPeriod)
.set("idempotencyKey", idempotencyKey)
.set("status", "pending")
.setOnInsert("createdAt", Instant.now()),
BillingAttempt.class
).flatMap(result ->
stripeClient.createPaymentIntent(stripeCustomerId, amountCents,
idempotencyKey))
.flatMap(paymentIntentId ->
ops.updateFirst(
Query.query(Criteria.where("idempotencyKey").is(idempotencyKey)),
Update.update("status", "succeeded")
.set("paymentIntentId", paymentIntentId),
BillingAttempt.class
).thenReturn(paymentIntentId))
).retryWhen(Retry.backoff(3, Duration.ofMillis(200))
.filter(e -> e instanceof StripeNetworkException));
}
The upsert pattern (find-or-create by idempotencyKey) replaces the raw save(). On attempt 1, the document is created. On any subsequent attempt with the same content-hash key, the upsert is a no-op on the document fields (the $setOnInsert for createdAt only applies on insert). The Stripe call proceeds with the same key — Stripe returns the already-committed result from its idempotency cache. The update to status: succeeded is idempotent as well. The entire operation is safe for any number of retries.
Mode 3: ReactiveMongoRepository.findAll() cold Flux + outer retryWhen() in batch billing — MongoDB re-queried on re-subscription, entire batch recharged
The developer is building a monthly batch billing job that iterates over all active customers in MongoDB, generates a Stripe charge per customer, and records the result. They use a ReactiveMongoRepository-derived query method to stream customers, apply flatMap to issue Stripe charges concurrently, and place retryWhen() on the outer Flux to handle transient infrastructure failures (network blips, MongoDB connection drops, Stripe rate-limit 429s).
The developer’s mental model of retryWhen() on a Flux: “If any element fails in the flatMap, retryWhen() retries from that element. MongoDB already returned all the previous customer records — they’re buffered or processed. The cursor is somewhere in the middle of the result set. Retry resumes from where the failure occurred.”
This mental model is incorrect about how Reactor’s retryWhen() works on a cold publisher. Reactor’s retryWhen() operator does not have any concept of “resume from position.” It re-subscribes the upstream publisher from the beginning. The upstream publisher here is a cold Flux<Customer> returned by the repository method. Re-subscription opens a new MongoDB cursor and re-queries from the start of the collection. All customers are re-emitted. Every flatMap invocation that generates a UUID and calls Stripe runs again for every customer.
// BatchBillingJob.java — unsafe Mode 3: findAll() cold Flux + outer retryWhen()
@Component
public class MongoReactiveBatchBillingJob {
private final CustomerRepository customerRepository;
private final BillingAttemptRepository billingAttemptRepository;
private final StripeWebClient stripeClient;
// Developer's reasoning:
// "customerRepository.findAllByActiveTrue() streams active customers from MongoDB.
// flatMap() sends a Stripe charge per customer (UUID generated inside flatMap — fresh
// idempotency key per customer per billing run). retryWhen() on the outer Flux handles
// transient failures. If a failure occurs mid-batch, retry resumes from the customer
// that failed — MongoDB already emitted the earlier customers."
//
// The problem:
// customerRepository.findAllByActiveTrue() returns a COLD Flux<Customer>.
// retryWhen() on the outer Flux re-subscribes the upstream Flux on any failure.
// Re-subscription triggers a new MongoDB find query. New cursor. All customers re-emitted.
// UUID.randomUUID() inside flatMap generates UUID_B for every customer. ← NEW UUIDs
// flatMap calls Stripe with UUID_B for all customers, including those charged in attempt 1.
// Blast radius = entire batch.
//
// Timeline (simplified — assume 3 customers, failure on customer 3):
// Subscription 1:
// MongoDB cursor opened. Customers C1, C2, C3 streamed.
// flatMap processes C1: UUID_A1 → Stripe → ch_A1 committed.
// flatMap processes C2: UUID_A2 → Stripe → ch_A2 committed.
// flatMap processes C3: UUID_A3 → Stripe → 500 returned. ← failure
// retryWhen() catches failure, schedules retry.
//
// Subscription 2 (retry):
// MongoDB re-queried. NEW cursor from beginning. C1, C2, C3 re-emitted.
// flatMap processes C1: UUID_B1 → Stripe → ch_B1. ← DUPLICATE for C1
// flatMap processes C2: UUID_B2 → Stripe → ch_B2. ← DUPLICATE for C2
// flatMap processes C3: UUID_B3 → Stripe → ch_B3. ← DUPLICATE for C3
public Mono<Void> runMonthlyBillingCycle(String billingPeriod) {
return customerRepository
.findAllByActiveTrue()
.flatMap(customer -> {
// UUID.randomUUID() inside flatMap — executes per customer per subscription.
// On retryWhen() re-subscription, this runs again for every re-emitted customer.
String idempotencyKey = UUID.randomUUID().toString();
return billingAttemptRepository
.save(new BillingAttempt(customer.getId(), customer.getAmountCents(),
billingPeriod, idempotencyKey))
.flatMap(attempt ->
stripeClient.createPaymentIntent(customer.getStripeCustomerId(),
customer.getAmountCents(), idempotencyKey));
}, 10) // concurrency of 10
.then()
.retryWhen(Retry.backoff(3, Duration.ofMillis(500))
.filter(e -> e instanceof TransientDataAccessException
|| e instanceof StripeNetworkException));
}
}
Why developers believe retry resumes mid-stream
The developer’s belief that retry resumes from a mid-stream position has an understandable origin. In imperative code, a for loop with a try/catch around the Stripe call does resume from the next element after the failure (or retries the specific element). The developer’s intuition is correct for that imperative pattern. The intuition fails when applied to a reactive retry operator on a cold source.
Reactor’s reactive streams specification (and the Reactive Streams specification in general) defines publishers as producers of zero-to-N elements. There is no protocol for “resume from position N.” A subscriber can cancel a subscription and re-subscribe, but re-subscription always starts from the beginning of a cold publisher. The retryWhen() operator is defined as: “re-subscribe the upstream publisher when the upstream emits an error signal, according to the companion publisher produced by the provided function.” There is no exception for cold repository query publishers — they are re-subscribed from the beginning like any other cold publisher.
The secondary misconception is about MongoDB cursor state. Once the error is emitted, the existing cursor is cancelled (via the Reactive Streams cancel() signal). The cursor position is not saved. A new subscription always opens a new cursor. There is no mechanism in the MongoDB reactive driver to resume an existing cursor by position — MongoDB cursors are server-side state associated with the connection/session that opened them. Cancelling and re-opening produces a fresh cursor from the beginning of the result set.
The fix: move retryWhen() inside flatMap for per-element retry scope
The fix moves retryWhen() from the outer Flux to inside each element’s flatMap lambda. Per-element retry re-subscribes only the element’s Mono on failure. The outer Flux<Customer> (the cold MongoDB query) is never re-subscribed. Customers who were already charged continue to the next element.
// BatchBillingJob.java — safe Mode 3 fix: per-element retry inside flatMap
@Component
public class MongoReactiveBatchBillingJob {
private final CustomerRepository customerRepository;
private final BillingAttemptRepository billingAttemptRepository;
private final StripeWebClient stripeClient;
public Mono<Void> runMonthlyBillingCycle(String billingPeriod) {
return customerRepository
.findAllByActiveTrue()
.flatMap(customer -> {
// Content-hash key: deterministic for this customer + period + amount.
// Computed before the per-element retry scope.
final String idempotencyKey = Hashing.sha256()
.hashString(customer.getId() + "|" + billingPeriod
+ "|" + customer.getAmountCents(), StandardCharsets.UTF_8)
.toString()
.substring(0, 36);
// Per-element Mono: save attempt + call Stripe.
Mono<String> chargeElement = billingAttemptRepository
.save(new BillingAttempt(customer.getId(),
customer.getAmountCents(), billingPeriod, idempotencyKey))
.flatMap(attempt ->
stripeClient.createPaymentIntent(
customer.getStripeCustomerId(),
customer.getAmountCents(), idempotencyKey));
// retryWhen() is INSIDE flatMap — applies to this customer's Mono only.
// On retry, only chargeElement is re-subscribed — idempotencyKey is
// a captured final variable and does not change. findAllByActiveTrue()
// is NOT re-subscribed.
return chargeElement
.retryWhen(Retry.backoff(3, Duration.ofMillis(200))
.filter(e -> e instanceof StripeNetworkException));
}, 10) // concurrency of 10
.then();
// No outer retryWhen() — prevents cold Flux re-subscription.
}
}
With retryWhen() inside flatMap, the retry scope is the individual element’s Mono. Re-subscribing chargeElement does not affect the outer Flux<Customer>. The MongoDB query is not re-executed. Customers C1 and C2, whose charges completed successfully before C3 failed, are not re-processed. C3’s chargeElement is re-subscribed with the same content-hash idempotency key — the captured idempotencyKey final variable does not change between retry attempts.
The content-hash key is especially important here. Even if the per-element chargeElement re-subscription triggers a new billingAttemptRepository.save() call (because the first save was part of a reactive pipeline that failed), the Stripe charge is deduplicated by key. If the MongoDB save succeeded on attempt 1 but Stripe returned a 500, the save on attempt 2 will either succeed (if the first save’s transaction rolled back) or throw a duplicate key error (if the first save committed). Handle both cases explicitly: use upsert by idempotency key for the MongoDB write (as shown in Mode 2’s fix), and use the same content-hash key for Stripe. This makes the per-element operation fully idempotent regardless of which sub-step failed.
Handling findAll() pagination for very large collections
For collections with hundreds of thousands of documents, findAllByActiveTrue() as a streaming query may hold a MongoDB server-side cursor open for a long time. If the job fails partway through and must restart, the entire collection is re-streamed from the beginning. The per-element retry handles transient Stripe errors within the run, but it does not help if the job is interrupted by a process crash or timeout.
A checkpoint-based pattern handles this: before starting the batch run, write a cursor checkpoint record to MongoDB containing the current runId and billingPeriod. For each successfully processed customer, upsert a ProcessedCustomer{runId, customerId, idempotencyKey, status} record. On restart, query for customers who do NOT have a ProcessedCustomer record for this runId. This is a MongoDB-idiomatic pattern using an anti-join query (or a $lookup with a null check), and it ensures the batch can resume without re-processing already-completed customers after a crash.
// Checkpoint-based pattern for resumable batch billing
public Mono<Void> runMonthlyBillingCycle(String runId, String billingPeriod) {
// Fetch customers who have NOT been processed in this run.
// Uses a MongoDB aggregation anti-join pattern.
return customerRepository
.findUnprocessedCustomersForRun(runId, billingPeriod) // custom query method
.flatMap(customer -> {
final String idempotencyKey = Hashing.sha256()
.hashString(customer.getId() + "|" + billingPeriod
+ "|" + customer.getAmountCents(), StandardCharsets.UTF_8)
.toString().substring(0, 36);
return billingAttemptRepository
.save(new BillingAttempt(customer.getId(),
customer.getAmountCents(), billingPeriod, idempotencyKey))
.flatMap(attempt ->
stripeClient.createPaymentIntent(
customer.getStripeCustomerId(),
customer.getAmountCents(), idempotencyKey))
.flatMap(paymentIntentId ->
processedCustomerRepository.save(
new ProcessedCustomer(runId, customer.getId(),
idempotencyKey, "succeeded", paymentIntentId)))
.retryWhen(Retry.backoff(3, Duration.ofMillis(200))
.filter(e -> e instanceof StripeNetworkException))
.onErrorResume(e -> {
// Log permanent failure; continue with remaining customers.
log.error("Permanent failure for customer {}: {}", customer.getId(), e.getMessage());
return processedCustomerRepository.save(
new ProcessedCustomer(runId, customer.getId(),
idempotencyKey, "failed", null))
.then(Mono.empty());
});
}, 10)
.then();
}
The findUnprocessedCustomersForRun custom query method uses a MongoDB aggregation that filters out customers already present in the processedCustomers collection for the given runId. On restart, the query returns only the remaining customers. The per-element idempotency key (content-hash of customerId + billingPeriod + amountCents) is stable across restarts — the same key will be sent to Stripe for any customer who was partially processed and is being retried in a subsequent job run.
Comparison: Spring Data MongoDB vs. Spring Data JPA reactive failure modes
| Dimension | Spring Data MongoDB | Spring Data JPA |
|---|---|---|
| AOP proxy ordering | Same: @Retryable outer (MAX-5), @Transactional inner (MAX) |
Same: @Retryable outer (MAX-5), @Transactional inner (MAX) |
| Transaction manager | ReactiveMongoTransactionManager — requires MongoDB 4.0+ replica set; manages ClientSession |
JpaTransactionManager — manages Hibernate EntityManager and JDBC connection |
| Transactional resource on rollback | MongoDB write rolled back (replica set oplog); Stripe HTTP call NOT rolled back | Hibernate EntityManager flush + JDBC rollback; Stripe HTTP call NOT rolled back |
| Audit gap on rollback | MongoDB audit record for UUID_A rolled back; only UUID_B record in MongoDB; Stripe shows both ch_A and ch_B | JPA entity for UUID_A rolled back (or not flushed); same audit gap pattern |
| Programmatic transaction API | ReactiveMongoTemplate.inTransaction() callback Function is re-invoked on retryWhen() re-subscription |
TransactionTemplate.execute() callback is re-invoked on proceed() (synchronous); TransactionalOperator.transactional() chain is re-assembled on retryWhen() re-subscription (reactive) |
| Cold publisher source | ReactiveMongoRepository.findAll() — backed by MongoDB reactive driver cursor; fresh cursor per subscription |
Spring Data JPA reactive does not exist; Spring Data R2DBC findAll() — backed by R2DBC cursor; similar cold semantics |
| Idempotent write fix | MongoDB upsert by idempotency key field with $setOnInsert |
JPA findOrCreateByIdempotencyKey with optimistic lock; or INSERT ... ON CONFLICT DO NOTHING via native query |
Detection signals for each failure mode
Each mode produces a characteristic observable signal in the Stripe dashboard and application logs that can be used to detect the issue in production before a customer complaint.
Mode 1 (@Retryable + @Transactional): Two PaymentIntent objects in the Stripe dashboard within seconds of each other for the same customer, same amount, same billing period — but with different Idempotency-Key values. The gap between ch_A.created and ch_B.created matches the @Backoff delay. In application logs, the transaction open/close log entries from ReactiveMongoTransactionManager show two transaction IDs (two ClientSession starts) within the same service operation span. Prometheus metric: duplicate (customerId, billingPeriod) pairs in the billing attempt collection within a short time window, with different idempotency keys.
Mode 2 (inTransaction() + retryWhen()): Same Stripe signal as Mode 1. In application logs, the inTransaction() debug output (if enabled via logging.level.org.springframework.data.mongodb.core=DEBUG) shows two transaction start/commit/rollback cycles within the same reactive subscription chain. The timing gap between ch_A and ch_B matches the Retry.backoff() delay. The MongoDB audit collection shows two BillingAttempt documents with the same customer + billing period but different idempotency keys (assuming the first transaction committed partially before Stripe failed).
Mode 3 (findAll() + outer retryWhen()): This is the most dangerous mode because the blast radius is the entire batch. All customers have two charges in Stripe within the retry window. Application logs show the findAll() query executing twice (visible in MongoDB driver debug logs as two find commands). The number of Stripe API calls doubles. The second set of charges has a uniform time offset matching the Retry.backoff() initial delay. Prometheus metric: total Stripe API calls in the billing window is approximately 2 × number_of_active_customers instead of 1×.
Test patterns for all three modes
All three failure modes can be verified with @DataMongoTest or @SpringBootTest using a WireMock HTTP server for Stripe and an embedded MongoDB via Flapdoodle or a Testcontainers MongoDB container.
// Mode 1 test — @Retryable + @Transactional UUID stability on Stripe transient failure
@SpringBootTest
@Testcontainers
class Mode1RetryableTransactionalTest {
@Container
static MongoDBContainer mongo = new MongoDBContainer("mongo:7.0")
.withCommand("--replSet", "rs0"); // Replica set required for transactions
@Autowired
private MongoReactiveBillingService billingService;
@Autowired
private BillingAttemptRepository billingAttemptRepository;
private WireMockServer wireMock;
@BeforeEach
void setUp() {
wireMock = new WireMockServer(wireMockConfig().dynamicPort());
wireMock.start();
// Configure the StripeWebClient to point at WireMock
// (via application property override or ReflectionTestUtils)
}
@Test
void retryableTransactional_usesStableIdempotencyKey_acrossRetryAttempts() {
// Stripe: first call returns 500, second call returns 200
wireMock.stubFor(post(urlPathEqualTo("/v1/payment_intents"))
.inScenario("retry")
.whenScenarioStateIs(STARTED)
.willReturn(aResponse().withStatus(500).withBody("{\"error\":{\"type\":\"api_error\"}}"))
.willSetStateTo("retried"));
wireMock.stubFor(post(urlPathEqualTo("/v1/payment_intents"))
.inScenario("retry")
.whenScenarioStateIs("retried")
.willReturn(aResponse().withStatus(200)
.withBody("{\"id\":\"pi_xxx\",\"status\":\"succeeded\"}")));
StepVerifier.create(billingService.chargeCustomer("C1", "cus_xxx", 5000L, "2026-10"))
.expectNextCount(1)
.verifyComplete();
// Assert: exactly TWO Stripe calls were made (one 500, one 200)
List<LoggedRequest> stripeRequests = wireMock.findAll(postRequestedFor(
urlPathEqualTo("/v1/payment_intents")));
assertThat(stripeRequests).hasSize(2);
// Assert: BOTH calls used the SAME Idempotency-Key header
String key1 = stripeRequests.get(0).getHeader("Idempotency-Key");
String key2 = stripeRequests.get(1).getHeader("Idempotency-Key");
assertThat(key1).isNotNull();
assertThat(key1).isEqualTo(key2); // ← This fails on the unsafe implementation
}
}
// Mode 2 test — inTransaction() callback UUID stability on retryWhen() re-subscription
@SpringBootTest
class Mode2InTransactionRetryTest {
@Autowired
private MongoProgrammaticBillingService billingService;
@Test
void inTransaction_withRetryWhen_usesStableIdempotencyKey() {
// WireMock setup: 500 then 200 (same as Mode 1)
// ... (same WireMock setup as above)
StepVerifier.create(billingService.chargeCustomer("C1", "cus_xxx", 5000L, "2026-10"))
.expectNextCount(1)
.verifyComplete();
List<LoggedRequest> requests = wireMock.findAll(postRequestedFor(
urlPathEqualTo("/v1/payment_intents")));
assertThat(requests).hasSize(2);
assertThat(requests.get(0).getHeader("Idempotency-Key"))
.isEqualTo(requests.get(1).getHeader("Idempotency-Key")); // fails on unsafe impl
}
}
// Mode 3 test — findAll() cold Flux: per-element retry does NOT re-query MongoDB
@SpringBootTest
@Testcontainers
class Mode3FindAllRetryTest {
@Container
static MongoDBContainer mongo = new MongoDBContainer("mongo:7.0");
@Autowired
private CustomerRepository customerRepository;
@Autowired
private MongoReactiveBatchBillingJob batchJob;
@Test
void batchBilling_withPerElementRetry_doesNotRechargeAlreadyProcessedCustomers() {
// Seed 3 customers
List<Customer> customers = List.of(
new Customer("C1", "cus_aaa", 1000L, true),
new Customer("C2", "cus_bbb", 2000L, true),
new Customer("C3", "cus_ccc", 3000L, true));
customerRepository.saveAll(customers).collectList().block();
// Stripe: succeed C1 and C2; fail C3 once then succeed
// (WireMock scenario with per-customer path matching)
wireMock.stubFor(post(urlPathEqualTo("/v1/payment_intents"))
.withRequestBody(containing("cus_aaa"))
.willReturn(okJson("{\"id\":\"pi_C1\",\"status\":\"succeeded\"}")));
wireMock.stubFor(post(urlPathEqualTo("/v1/payment_intents"))
.withRequestBody(containing("cus_bbb"))
.willReturn(okJson("{\"id\":\"pi_C2\",\"status\":\"succeeded\"}")));
wireMock.stubFor(post(urlPathEqualTo("/v1/payment_intents"))
.withRequestBody(containing("cus_ccc"))
.inScenario("c3-retry").whenScenarioStateIs(STARTED)
.willReturn(aResponse().withStatus(500))
.willSetStateTo("c3-retried"));
wireMock.stubFor(post(urlPathEqualTo("/v1/payment_intents"))
.withRequestBody(containing("cus_ccc"))
.inScenario("c3-retry").whenScenarioStateIs("c3-retried")
.willReturn(okJson("{\"id\":\"pi_C3\",\"status\":\"succeeded\"}")));
StepVerifier.create(batchJob.runMonthlyBillingCycle("2026-10"))
.verifyComplete();
// Assert: C1 and C2 each called exactly once (not re-called on C3's retry)
List<LoggedRequest> c1Requests = wireMock.findAll(postRequestedFor(
urlPathEqualTo("/v1/payment_intents"))
.withRequestBody(containing("cus_aaa")));
List<LoggedRequest> c2Requests = wireMock.findAll(postRequestedFor(
urlPathEqualTo("/v1/payment_intents"))
.withRequestBody(containing("cus_bbb")));
List<LoggedRequest> c3Requests = wireMock.findAll(postRequestedFor(
urlPathEqualTo("/v1/payment_intents"))
.withRequestBody(containing("cus_ccc")));
assertThat(c1Requests).hasSize(1); // ← fails on unsafe impl (hasSize(2))
assertThat(c2Requests).hasSize(1); // ← fails on unsafe impl (hasSize(2))
assertThat(c3Requests).hasSize(2); // 1 failure + 1 success (correct)
// Assert: C3's two requests used the same idempotency key (per-element retry)
assertThat(c3Requests.get(0).getHeader("Idempotency-Key"))
.isEqualTo(c3Requests.get(1).getHeader("Idempotency-Key"));
}
}
The Mode 3 test is the most important for production confidence because it verifies the assertion most likely to be wrong: that C1 and C2 were each called exactly once. The unsafe implementation (outer retryWhen()) causes findAll() to re-query, which causes C1 and C2 to be re-emitted and re-charged. assertThat(c1Requests).hasSize(1) catches this definitively.
Summary: the mental model for Spring Data MongoDB + Stripe idempotency
All three failure modes share a common root cause: UUID generation is inside a code boundary that is re-executed by the retry mechanism. The mechanisms differ — AOP proceed() in Mode 1, reactive callback re-invocation in Mode 2, cold publisher re-subscription in Mode 3 — but the structural fix is always the same:
- Place UUID computation above (outside) the retry boundary.
- Pass UUID as a captured final variable or stable method parameter into the retried code.
- Use a content-hash key derived from stable billing inputs when possible, so that independent concurrent invocations for the same customer + period + amount naturally converge to the same key.
For Spring Data MongoDB specifically, the three boundaries to be aware of are:
@Retryableon a method with@Transactional: the resume point is the beginning of the method body (line 1), not the Stripe call. UUID at method entry is inside the retry boundary.ReactiveMongoTemplate.inTransaction()callbackFunction: whenretryWhen()is applied to the outerMonoreturned byinTransaction(), the callbackFunctionis re-invoked. Any deferred computation inside the callback body re-evaluates. UUID computed withMono.defer{}orMono.fromCallable{}inside the callback is inside the retry boundary. UUID computed as afinallocal variable in the enclosing method (before theinTransaction()call) is outside it.ReactiveMongoRepositoryquery methods returning coldFlux:retryWhen()applied to the outerFluxre-subscribes from the MongoDB cursor start. The entire document stream is re-emitted. UUID generated insideflatMapruns again for every re-emitted document. MoveretryWhen()insideflatMapfor per-element retry scope.
A proxy that sits between your agent and Stripe — one that tracks the idempotency key sent per logical billing operation and alerts when two different keys are used for the same customer, same amount, and same billing period within a short window — would catch all three of these failure modes in production. That is exactly what Keybrake is building: a scoped API-key proxy for the non-LLM SaaS APIs your agent calls, with per-vendor spend caps, audit log, and duplicate-charge detection at the proxy layer.
Protect your agent’s Stripe keys
Keybrake proxies the Stripe, Twilio, and Resend API calls your agent makes — enforcing per-day spend caps, logging every charge with the idempotency key used, and flagging duplicate-key anomalies before they become duplicate charges. Early access waitlist below.