Project Reactor并行读取场景下如何实现数据缓冲批量写入
Reactor 流式攒批写入数据库实现方案
核心思路
直接使用Reactor内置的攒批操作符即可实现需求,响应式流规范原生支持背压,无需手动实现读取暂停逻辑:
- 用
buffer操作符将逐条发送的消息组装为指定大小的批次集合 - 下游写入速度不足时,Reactor会自动降低上游数据请求频率,数据源会自动暂停拉取新消息
修正后的完整代码
首先修正原有代码的类型不匹配、类命名不规范问题:
// 类名遵循Java大驼峰规范 class DataSource{ public Item getNextItem(){ // 原有消息拉取逻辑 } } class DataTarget{ public void write(List<Item> items){ // 原有批量写入逻辑 } }
主流程实现代码:
Flux.generate((SynchronousSink<Item> sink) -> { try { Item val = dataSource.getNextItem(); if (val == null) { sink.complete(); return; } sink.next(val); } catch (InterruptedException e) { sink.error(e); } }) .parallel(4) // IO密集型写入用boundedElastic调度器更合适 .runOn(Schedulers.boundedElastic()) // 每10条攒为一个批次,流结束时剩余不足10条也会自动输出 .buffer(10) // 批量写入 .doOnNext(dataTarget::write) .sequential() .blockLast();
关键配置说明
- 自定义批次大小:修改
buffer(10)的数值即可,建议根据数据库批量写入性能测试调整,通常50~200条为最优区间。 - 超时攒批:如果需要避免少量数据长时间滞留缓冲区,可使用
bufferTimeout(10, Duration.ofSeconds(1))替换buffer(10),表示满10条或等待1秒两个条件满足任意一个就触发写入。 - 并行度调整:修改
parallel(4)的数值调整并行消费和写入的线程数,可根据消息队列的消费速度、数据库的写入负载灵活调整。 - 错误处理:可根据业务需求添加重试、降级逻辑,比如追加
retry(3)实现写入失败自动重试3次,或追加onErrorContinue((throwable, batch) -> log.error("批次写入失败", throwable))跳过失败批次避免整条流中断。
内容的提问来源于stack exchange,提问作者Roman
相关产品推荐
相关产品推荐

