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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 01:45:30