RxPY并发操作串行化原因及线程亲和并行处理实现方案
现象原因
核心是RxPY的观察者默认带串行通知保护:
- 你在
flat_map下游写的map、subscribe逻辑,属于同一个下游观察者实例,RxPY默认会给这个观察者套一层SerializedObserver包装,保证同一时间永远只有一个on_next调用在执行,不允许并发重入,避免流处理顺序混乱、线程安全问题。 - 你的代码里
flat_map确实把5个初始sleep任务并行提交给了线程池,5个任务会分别在1、2、3、4、5秒的时间点完成,但完成后要往下游发数据时,必须抢下游观察者的锁:第一个任务1秒时抢到锁,执行map里的2秒sleep、再执行subscribe逻辑,总共占锁3秒才释放;第二个任务哪怕2秒就跑完了,也得等锁释放,3秒时才能抢到锁进map,以此类推,最终就出现了map触发时间依次间隔2秒的排队现象。 - 你看到不同元素用了不同线程,只是因为各个future完成的线程不一样,抢锁成功的线程不同而已,本质还是串行排队执行。
目标效果实现方案
要实现「5个线程并行处理、每个元素全流程同线程、无额外排队开销」的效果,核心是不要把元素处理逻辑放在flat_map下游的公共串行管道里,而是把每个元素从初始sleep到map、subscribe的全流程逻辑,都作为独立任务提交给线程池,让每个任务在自己分配的线程上独立跑完,避开下游公共观察者的串行锁。
实现代码如下:
import reactivex as rx import concurrent.futures import time from reactivex import operators as ops from threading import current_thread # 单元素全流程处理逻辑,全程在同一个线程执行 def handle_item(x, program_start): # 原map逻辑 print('map', current_thread().name, time.time() - program_start, x) time.sleep(2) # 原subscribe逻辑 print('sub', current_thread().name, time.time() - program_start, x) return x if __name__ == "__main__": with concurrent.futures.ThreadPoolExecutor(max_workers=5) as executor: start_ts = time.time() # 阻塞等待所有元素处理完成再退出,避免线程池提前关闭 rx.from_iterable(range(1, 6)).pipe( ops.flat_map(lambda val: rx.from_future(executor.submit(handle_item, val, start_ts))), ops.to_list() ).run()
运行后输出逻辑和你给出的pseq效果完全一致:
- 1秒时值为1的任务完成初始sleep,打印map,再睡2秒后打印sub
- 2秒时值为2的任务完成初始sleep,打印map,再睡2秒后打印sub
- 3秒时值为1的任务打印sub,同时值为3的任务完成初始sleep打印map
- 后续依次执行,总耗时约7秒,无额外排队等待
- 每个元素的map、sub打印都是同一个线程,满足线程亲和性要求。
内容的提问来源于stack exchange,提问作者Shuming
相关产品推荐
相关产品推荐

