Spring Webflux非阻塞操作:无block()从Flux取值并联动查询
Hey there! Let's break down how to implement your requirements in a fully non-blocking way, and fix the gaps in your current code snippet.
First, Let's Address the Issues in Your Current Approach
Your existing code has a couple of key problems that prevent it from meeting your non-blocking goals:
- When you use
Flux.just(product.getAvailabilityCalendar()), you're wrapping the entireArrayList<ProductAvailability>as a single element in the Flux. This doesn't let you process individualProductAvailabilityitems for your third query step. - You haven't connected the extracted availability data to the subsequent query operation, which is critical to completing your end-to-end workflow.
Correct Non-Blocking Implementation
Here's a revised version of your method that follows all your requirements, with clear explanations for each step:
// Adjust the return type if you need to return the final queried data instead of Product public Flux<Product> getAllProductsByAvailability(Flux<ProductProperties> productProperties, Map<String, String> searchParams) { return productProperties // Step 1: Fetch all Products for each ProductProperties (non-blocking via flatMap) .flatMap(property -> productRepository.findByProductPropertiesId(property.getId())) // Step 2: Extract individual ProductAvailability items from each Product's calendar // Flux.fromIterable converts the ArrayList into a stream of elements (fully non-blocking) .flatMapMany(product -> Flux.fromIterable(product.getAvailabilityCalendar()) // Step 3: Use each ProductAvailability to query additional data (replace with your actual query method) .flatMap(availability -> yourOtherRepository.findByAvailabilityId(availability.getId())) // Optional: If you need to map back to the original Product (or combine with queried data) .map(queriedData -> combineProductAndRelatedData(product, queriedData)) ); } // Example helper method if you need to combine Product with the queried data private Product combineProductAndRelatedData(Product product, QueriedDataType queriedData) { // Update the product with the queried data (e.g., attach related info to the product) return product; }
Key Non-Blocking Principles Followed
- No
block()calls: We use reactive operators likeflatMap,flatMapMany, andFlux.fromIterableto handle all operations without blocking threads—this keeps the reactive pipeline flowing as intended. - Stream processing: Each step processes elements as they become available, rather than waiting for the entire Flux to complete, which maximizes efficiency.
- Proper 1-to-many transformation:
flatMapManyis used here because we're converting a singleProductinto multipleProductAvailabilityelements, which aligns perfectly with reactive stream patterns.
Notes on Return Type
If your final goal is to return the data from the third query step (instead of the original Product), simply adjust the method's return type to match the type returned by yourOtherRepository.findByAvailabilityId() (e.g., Flux<QueriedDataType>).
Is Your Original Approach Feasible?
Unfortunately, no. The original code doesn't properly split the availabilityCalendar into individual elements, and it fails to link to the third query step. By adjusting the operators to correctly transform the Flux streams, you can achieve a fully non-blocking workflow as intended.
内容的提问来源于stack exchange,提问作者Matexon

