
当使用 buffer() 分批数据后,需通过 parallel() + runOn() 而非 publishOn() 实现真正的多线程并行处理;否则 map 仍串行运行于单一线程。
当使用 `buffer()` 分批数据后,需通过 `parallel()` + `runon()` 而非 `publishon()` 实现真正的多线程并行处理;否则 `map` 仍串行运行于单一线程。
在 Project Reactor 中,Flux 默认是顺序、单线程流水线式执行的。即使你调用 .publishOn(Schedulers.parallel()),它仅影响下游操作符的调度上下文(即后续 doOnNext、flatMap 等在指定线程执行),但 map 本身仍由上游 Publisher 在当前线程逐个触发——尤其当上游是冷流(如 Flux.just)且数据瞬间就绪时,buffer(2) 生成的 List<string></string> 会一次性全部发出,随后 map 在同一个线程上连续执行,无法并发。
✅ 正确做法是启用 并行化处理模型:
使用 .parallel() 将 Flux<list>></list> 转换为 ParallelFlux(逻辑分片),再通过 .runOn(Schedulers.parallel()) 为每个分片分配独立线程执行 map,最后用 .sequential() 合并结果为有序 Flux:
@Test
public void testParallelMapOnBufferedFlux() throws InterruptedException {
Flux.just("a", "b", "c", "d", "e", "f", "g")
.buffer(2)
.parallel() // ✅ 启用并行处理(默认并行度 = CPU 核心数)
.runOn(Schedulers.parallel()) // ✅ 每个分片在独立 parallel 线程中执行
.map(this::doTimeintensiveStuff)
.sequential() // ✅ 恢复为普通 Flux(保持原始顺序)
.doOnNext(val -> {
log.info("[" + Thread.currentThread().getName() + "] " + val);
})
.blockLast();
}
private String doTimeintensiveStuff(List<string> input) {
try {
TimeUnit.MILLISECONDS.sleep(100); // 模拟耗时操作
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
return String.join(", ", input);
}</string>
? 关键区别说明:
-
publishOn(scheduler):切换订阅者线程,仅改变后续操作符的执行线程,不改变执行模型(仍是串行 pull); -
parallel()+runOn(scheduler):显式启用并行流水线,将数据分片并行调度到多个线程,map等中间操作在各自线程中真正并发执行。
⚠️ 注意事项:
Orderly React SDK 钩子使用参考指南,包括 useOrderEntry、usePositionStream、useOrderbookStream、useCollateral 等。
-
parallel()默认并行度为Runtime.getRuntime().availableProcessors(),可通过.parallel(n)自定义; - 若
map内部有共享状态或非线程安全操作(如写入静态变量、修改全局集合),必须自行加锁或使用线程安全结构; -
sequential()保证输出顺序与输入分片顺序一致(即["a","b"]→["c","d"]→ ...),但不保证各map完成时间顺序; - 避免在
map中执行阻塞 I/O;如需异步非阻塞调用,请改用flatMap+Mono.fromCallable(...).subscribeOn(...)。
掌握 parallel()/runOn()/sequential() 组合,是高效利用多核资源处理批量数据的关键。










