如何在Flux上实现自然批处理(智能批处理)?相关问题咨询
在Flux上实现自然批处理的方案
什么是自然批处理
自然批处理是基于单线程无限循环的批处理模式,核心流程为:
- 收集队列中所有待处理项(至少包含1个)
- 对这批项执行批量操作(如批量插入数据库)
和Project Reactor常规批处理逻辑不同,它不需要用take(x)限定批次大小,而是要在收集下一批数据的同时处理上一批,还要限制预取下一批的项数。
一、简便实现:组合现有Reactor操作符
Reactor原生没有直接对应自然批处理的操作符,但可以通过bufferUntil、publishOn和背压控制的组合来近似实现需求:
Flux<Integer> source = Flux.range(1, 1000) .onBackpressureBuffer(10); // 限制上游预取的缓存大小 source.bufferUntil(item -> false, true) // 收集当前所有可用元素作为一个批次 .publishOn(Schedulers.single(), 1) // 单线程处理批次,仅预取1个批次 .subscribe(batch -> { // 执行批量业务操作,比如批量写入数据库 System.out.println("处理批次,元素数量:" + batch.size()); });
bufferUntil(item -> false, true):第二个参数设为true时,会一次性收集当前流中所有可用元素组成批次;publishOn(Schedulers.single(), 1):指定单线程处理批次,同时限制预取1个批次,确保处理上一批时上游仅预取限定数量的下一批元素;onBackpressureBuffer(10):控制上游预取的元素上限,避免无限制缓存。
二、精准实现:自定义Publisher
自定义Publisher是完全可行的,且能精准贴合自然批处理的核心逻辑。通过自定义实现Reactive Streams规范的Publisher,我们可以直接控制批次收集、处理和预取的逻辑:
自定义Publisher示例
public class NaturalBatchPublisher<T> implements Publisher<List<T>> { private final Publisher<T> source; private final int prefetchLimit; public NaturalBatchPublisher(Publisher<T> source, int prefetchLimit) { this.source = source; this.prefetchLimit = prefetchLimit; } @Override public void subscribe(Subscriber<? super List<T>> subscriber) { new NaturalBatchSubscriber<>(subscriber, prefetchLimit).subscribe(source); } private static class NaturalBatchSubscriber<T> extends BaseSubscriber<T> { private final Subscriber<? super List<T>> downstream; private final int prefetchLimit; private final Queue<T> queue = new ConcurrentLinkedQueue<>(); private boolean sourceDone; public NaturalBatchSubscriber(Subscriber<? super List<T>> downstream, int prefetchLimit) { this.downstream = downstream; this.prefetchLimit = prefetchLimit; } @Override protected void hookOnSubscribe(Subscription subscription) { downstream.onSubscribe(new Subscription() { @Override public void request(long n) { // 下游请求n个批次,对应向上游请求n*prefetchLimit个元素 request(n * prefetchLimit); } @Override public void cancel() { NaturalBatchSubscriber.this.cancel(); } }); // 初始请求限定数量的元素 request(prefetchLimit); } @Override protected void hookOnNext(T value) { queue.add(value); emitBatch(); } @Override protected void hookOnComplete() { sourceDone = true; emitBatch(); } @Override protected void hookOnError(Throwable throwable) { downstream.onError(throwable); } private void emitBatch() { // 单线程循环收集所有可用元素组成批次 while (!queue.isEmpty()) { List<T> batch = new ArrayList<>(); T item; while ((item = queue.poll()) != null) { batch.add(item); } downstream.onNext(batch); // 根据当前队列剩余量,补充请求元素以维持预取上限 long needRequest = prefetchLimit - queue.size(); if (needRequest > 0) { request(needRequest); } } // 上游完成且队列空时,通知下游完成 if (sourceDone && queue.isEmpty()) { downstream.onComplete(); } } } }
使用自定义Publisher
Flux<Integer> dataSource = Flux.range(1, 1000); new NaturalBatchPublisher<>(dataSource, 10) .subscribe(batch -> { // 执行批量操作 System.out.println("处理批次,元素数量:" + batch.size()); });
这个实现完全遵循自然批处理的逻辑:单线程循环收集所有可用元素成批次,处理批次的同时动态控制上游预取的元素数量,完美匹配需求。
内容的提问来源于stack exchange,提问作者Lubo
相关产品推荐
相关产品推荐

