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

RxJava技术问题:如何实现生产者等待订阅者完成批量流处理?

Solution for Controlled Producer-Consumer Batching with RxJava

Let's break down your problem and fix the memory issue while keeping the producer able to pre-prepare the next batch.

Why Your Current Code Fails

  • When you add observeOn(Schedulers.io()): The observeOn operator queues all upstream events on the IO scheduler's unbounded queue. Combined with flatMap's default unbounded concurrency, your producer will keep generating objects non-stop, filling up the queue until you hit an OutOfMemoryError.
  • When you remove observeOn: The entire pipeline runs synchronously. The producer waits for the consumer to finish processing a batch before generating any new objects—so you lose the ability to pre-prepare the next batch.

The Fix: Controlled Concurrency + Backpressure

We need to let the producer run in parallel to pre-generate the next batch, but limit how many batches are being processed at once (so we don't overload memory). Here's the adjusted code:

Set<Bar> bars = Flowable.fromStream(foos)
    .subscribeOn(Schedulers.io()) // Run producer on a separate thread to pre-generate
    .map(foo -> calculateBars(foo))
    .buffer(10) // Batch into groups of 10
    // Limit to 1 concurrent batch processing, run save on IO thread
    .flatMap(
        batch -> Flowable.fromCallable(() -> saveBars(batch))
            .subscribeOn(Schedulers.io()),
        1 // Critical: Only process 1 batch at a time
    )
    .collect(Collectors.toSet())
    .blockingGet();

How This Works

  1. Parallel Producer: subscribeOn(Schedulers.io()) makes your producer run on an IO thread, so it can start generating the next batch of 10 objects while the current batch is being saved to the database.
  2. Batching: buffer(10) groups objects into manageable batches.
  3. Controlled Concurrency: The flatMap overload with maxConcurrency=1 ensures only one batch is processed at a time. RxJava's backpressure mechanism will automatically pause the producer once there's one unprocessed batch waiting—preventing unlimited memory accumulation.
  4. Non-Blocking Saving: Wrapping saveBars in Flowable.fromCallable() + subscribeOn(Schedulers.io()) ensures the database save operation doesn't block the producer thread.

Key Behavior

  • The producer will generate up to 2 batches total in memory at any time: one being processed by the consumer, and one pre-generated in the buffer waiting to be processed. This keeps memory usage constant.
  • As soon as the consumer finishes a batch, it immediately picks up the pre-generated next batch, and the producer starts generating the following batch right away.

内容的提问来源于stack exchange,提问作者Lutosław

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.27 13:37:29