Project Reactor分组求和时reduce调用block出现挂起问题如何解决
Project Reactor嵌套block导致流程挂起的解决方案
问题根因
你遇到的流程挂起是嵌套阻塞调用导致的死锁:
collectMap的value转换函数运行在Reactor流处理线程上,你在该函数内调用group.reduce(...).block()会直接占用当前线程等待reduce结果- 而
group.reduce需要上游的GroupedFlux持续发射数据才能完成计算,上游的执行线程已经被你主动block占用,最终形成死锁导致流程卡死。
正确实现方案
不要在collectMap的转换函数中嵌套block,先对每个分组单独完成求和计算,再统一收集为Map即可,同时注意两个阻塞采集逻辑的调度器配置,避免阻塞非弹性线程。
修正后完整代码如下:
import org.junit.jupiter.api.Test; import reactor.core.publisher.Flux; import reactor.core.publisher.FluxSink; import reactor.core.scheduler.Schedulers; import java.math.BigDecimal; import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.AbstractMap; public class TestBlockingIssue { @Test public void testBlockingMap() { // 阻塞的通用数据采集逻辑放到boundedElastic调度器,避免阻塞主线程 final Flux<Wrapper> source = Flux.create( sink -> createWrappers( 1000, sink ) ) .subscribeOn(Schedulers.boundedElastic()); final Map<String, BigDecimal> result = source.parallel( 10 ).runOn( Schedulers.boundedElastic() ) .map( wrapper -> enhanceWrapper( wrapper, 100 ) ) .flatMap( wrapper -> Flux.fromIterable( wrapper.detailsList ) ) .sequential() .groupBy( details -> details.detailKey ) // 先对每个分组单独求和,再拼成键值对 .flatMap(group -> group.reduce(BigDecimal.ZERO, (acc, detail) -> acc.add(detail.value)) // JDK16以下可替换为 new AbstractMap.SimpleEntry<>(group.key(), sum) .map(sum -> Map.entry(group.key(), sum))) // 最后统一收集为Map,无嵌套阻塞 .collectMap(Map.Entry::getKey, Map.Entry::getValue) .block(); System.out.println( result ); } private Wrapper enhanceWrapper( final Wrapper wrapper, final int count ) { for ( int i = 0; i < count; i++ ) { wrapper.detailsList.add( new Details( "detail-" + i, new BigDecimal( i +1 ) ) ); } return wrapper; } private void createWrappers( final int count, final FluxSink<Wrapper> sink ) { for ( int i = 0; i < count; i++ ) { sink.next( new Wrapper( "Wrapper-" + i ) ); } sink.complete(); } private class Details { final String detailKey; final BigDecimal value; private Details( final String detailKey, final BigDecimal value ) { this.detailKey = detailKey; this.value = value; } } private class Wrapper { final String lookupKey; final List<Details> detailsList = new ArrayList<>(); private Wrapper( final String lookupKey ) { this.lookupKey = lookupKey; } } }
运行后输出的Map完全符合你的预期:detail-0对应值1000、detail-1对应值2000,以此类推。
内容的提问来源于stack exchange,提问作者sarfog
相关产品推荐
相关产品推荐

