解决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)
问题原因
- 嵌套订阅破坏背压:在
onNext中直接调用subscribe()启动saveAll流,这种嵌套订阅会脱离Reactor的背压机制,导致上游(KafkaReceiver + bufferTimeout)无法感知下游的处理速度,当saveAll处理较慢时,bufferTimeout的缓冲区会不断积压数据,最终触发溢出。 - 不合理的buffer参数:使用
Integer.MAX_VALUE作为批次最大容量,会导致buffer一次性缓存大量数据,加剧内存压力和溢出风险。 - 手动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
相关产品推荐
相关产品推荐

