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

解决Reactor OverflowException:Spring Data响应式仓库保存Flux报错

解决Spring Data响应式仓库保存Flux数据时的OverflowException异常

需求描述

使用Spring Data响应式仓库(Elasticsearch、MongoDB、Cassandra均可)保存Kafka消费到的Flux数据流。

尝试的实现代码

public Flux<Foo> issue() {
        Sinks.Many<Foo> sink = Sinks.many().unicast().onBackpressureBuffer();
        kafkaReceiver.receive()
                .map(message -> convertMesageToFoo(message)
                )
                .bufferTimeout(Integer.MAX_VALUE, Duration.ofMillis(500))//尝试过多种maxSize和maxTime参数
                .subscribe(new Subscriber<>() {
                    private Subscription subscription;
                    volatile boolean upstreamComplete = false;

                    @Override
                    public void onSubscribe(Subscription subscription) {
                        this.subscription = subscription;
                        subscription.request(1);
                    }

                    @Override
                    public void onNext(List<Foo> foos) {
                        reactiveElasticRepository.saveAll(foos)
                                .map(sink::tryEmitNext)
                                .doOnComplete(() -> {
                                    if (!upstreamComplete) {
                                        subscription.request(1);
                                    } else {
                                        sink.tryEmitComplete();
                                    }
                                })
                                .subscribe();
                    }

                    @Override
                    public void onError(Throwable throwable) {
                        throwable.printStackTrace();
                        System.err.println("该问题可100%复现:" + throwable.getLocalizedMessage());
                        subscription.cancel();
                        sink.tryEmitError(throwable);
                    }

                    @Override
                    public void onComplete() {
                        upstreamComplete = true;
                    }
                });

        return sink.asFlux();
    }

遇到的问题

该问题可100%复现,抛出如下异常:

reactor.core.Exceptions$ErrorCallbackNotImplemented: reactor.core.Exceptions$OverflowException: Could not emit buffer due to lack of requests
Caused by: reactor.core.Exceptions$OverflowException: Could not emit buffer due to lack of requests
    at reactor.core.Exceptions.failWithOverflow(Exceptions.java:249)
    at reactor.core.publisher.FluxBufferTimeout$BufferTimeoutSubscriber.flushCallback(FluxBufferTimeout.java:227)
    at reactor.core.publisher.FluxBufferTimeout$BufferTimeoutSubscriber.lambda$new$0(FluxBufferTimeout.java:158)
    at reactor.core.scheduler.WorkerTask.call(WorkerTask.java:84)
    at reactor.core.scheduler.WorkerTask.call(WorkerTask.java:37)
    at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
    at java.base/java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:304)
    at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
    at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
    at java.base/java.lang.Thread.run(Thread.java:833)

问题原因

  1. 嵌套订阅破坏背压:在onNext中直接调用subscribe()启动saveAll流,这种嵌套订阅会脱离Reactor的背压机制,导致上游(KafkaReceiver + bufferTimeout)无法感知下游的处理速度,当saveAll处理较慢时,bufferTimeout的缓冲区会不断积压数据,最终触发溢出。
  2. 不合理的buffer参数:使用Integer.MAX_VALUE作为批次最大容量,会导致buffer一次性缓存大量数据,加剧内存压力和溢出风险。
  3. 手动Subscriber的背压控制不严谨:手动实现的Subscriber在处理saveAll完成后才请求下一批,但bufferTimeout的定时flush逻辑不受这个控制,当定时触发时如果没有足够的请求许可,就会抛出OverflowException。

解决方案

修正后的代码

public Flux<Foo> issue() {
    return kafkaReceiver.receive()
            // 保留Kafka消息用于后续提交偏移量,同时转换为Foo对象
            .map(message -> new AbstractMap.SimpleEntry<>(message, convertMesageToFoo(message)))
            // 合理设置批次大小和超时,平衡写入效率与内存占用
            .bufferTimeout(1000, Duration.ofMillis(500))
            // 顺序处理每个批次,确保背压在整个链路中正确传递
            .concatMap(batch -> {
                // 提取批次中的Foo对象
                List<Foo> foos = batch.stream()
                        .map(AbstractMap.SimpleEntry::getValue)
                        .collect(Collectors.toList());
                // 批量保存到Elasticsearch
                return reactiveElasticRepository.saveAll(foos)
                        // 保存完成后批量提交Kafka偏移量,确保数据一致性
                        .doOnComplete(() -> batch.forEach(entry -> entry.getKey().acknowledge()));
            })
            .doOnError(throwable -> {
                throwable.printStackTrace();
                System.err.println("异常信息:" + throwable.getLocalizedMessage());
            });
}

关键改进点

  • 移除嵌套订阅:用concatMap替代手动Subscriber,concatMap会顺序处理每个批次的saveAll流,只有当前批次处理完成后才会请求下一批数据,确保背压在整个链路中正确传递。
  • 合理设置批次参数:将bufferTimeout的maxSize调整为合理值(如1000),避免缓存过多数据导致内存溢出。
  • 移除冗余Sink:直接返回链式处理后的Flux,无需额外通过Sink转发数据,减少复杂度和潜在的背压问题。
  • 批量处理Kafka偏移量:在批次保存完成后统一提交偏移量,避免消息丢失或重复消费。

可选优化

如果需要并行处理多个批次以提高写入效率,可以将concatMap替换为flatMap并指定并发数(需确保数据库能承受对应并发压力):

.flatMap(batch -> {
    // 同concatMap中的处理逻辑
}, 3) // 允许同时处理3个批次

内容的提问来源于stack exchange,提问作者PatPanda

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 00:00:03