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()): TheobserveOnoperator queues all upstream events on the IO scheduler's unbounded queue. Combined withflatMap's default unbounded concurrency, your producer will keep generating objects non-stop, filling up the queue until you hit anOutOfMemoryError. - 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
- 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. - Batching:
buffer(10)groups objects into manageable batches. - Controlled Concurrency: The
flatMapoverload withmaxConcurrency=1ensures 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. - Non-Blocking Saving: Wrapping
saveBarsinFlowable.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
相关产品推荐
相关产品推荐

