基于Spring-WebFlux实现无阻塞双集合插入并回滚方案咨询
Optimizing Reactor-Based Two-Step Insert with Rollback (Non-Blocking, 2 DB Calls)
Great question—your current implementation does have a couple of key issues (manual subscribe() calls breaking the reactive chain, unreliable rollback handling) that we can fix while keeping the minimal 2 database calls for the happy path. Let's walk through the solution step by step.
Problems with the Current Code
First, let's call out what's not working optimally:
- Blocking Risk: Using
subscribe()insidedoOnSuccesskicks off the details insert as a separate, unmanaged asynchronous operation. This breaks the reactive chain, meaning the outerMono<Business>might complete before the details insert finishes, and you lose visibility into its success/failure. - Unreliable Rollback: The manual
businessRepository.delete(businezz).subscribe()in the error handler isn't guaranteed to execute properly (e.g., if the error happens before the subscription completes), and it doesn't propagate the error back to the caller. - Non-Chained Logic: The flow isn't a single reactive pipeline, so you can't leverage Reactor's built-in error handling and backpressure.
The Solution: Reactive Chaining with flatMap and onErrorResume
We can rewrite this using proper reactive operators to keep everything non-blocking, maintain the 2-database-call happy path, and ensure atomic rollback on failure. Here's the optimized Kotlin code:
override fun insert(business: Business): Mono<Business> { return businessRepository.save(business) .flatMap { savedBusiness -> // Prepare the details with the saved business ID val businessDetails = savedBusiness.businessDetails!!.apply { idBusiness = savedBusiness.id } // Chain the details insert, then return the original saved business on success businessDetailsService.insert(businessDetails) .thenReturn(savedBusiness) // Handle rollback if details insert fails .onErrorResume { error -> businessRepository.delete(savedBusiness) .then(Mono.error(error)) // Propagate the original error after rollback } } }
How This Works
Let's break down the key parts:
flatMapfor Chained Asynchronous Operations: Instead ofdoOnSuccess, we useflatMapto embed the details insert directly into the reactive pipeline. This keeps all operations non-blocking and ensures the outerMonoonly completes when both steps are done (or fails if either step fails).thenReturn(savedBusiness): After the details insert succeeds, we return the original saved business—exactly what you need (return the first operation's result when both succeed).onErrorResumefor Atomic Rollback: If the details insert throws an error, we first execute thedeleteoperation to rollback the business insert, then rethrow the original error so the caller knows the operation failed. This guarantees that if the second step fails, the first step's changes are undone.- Minimal Database Calls: The happy path only uses 2 database operations (
save(business)andinsert(businessDetails)). The rollbackdeleteis only executed if there's an error, which is unavoidable for consistency.
Why This Is Better
- Fully Non-Blocking: No manual
subscribe()calls mean the entire pipeline respects Reactor's reactive model, with proper backpressure and asynchronous execution. - Reliable Error Handling: Rollback is part of the reactive chain, so it's guaranteed to run if the second operation fails.
- Clear, Maintainable Flow: The code reads as a single logical pipeline, making it easier to debug and extend later.
内容的提问来源于stack exchange,提问作者David Latief Budiman
相关产品推荐
相关产品推荐

