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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 05:24:07