You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.06 02:27:01