RXPY v3如何实现批处理串行执行?替代自定义信号量方案
使用RXPY v3实现批次串行、批内并行的任务执行
你不需要用全局变量模拟信号量,RXPY v3提供了**exhaustMap**运算符,完美匹配你的需求:上一批任务全部结束后再启动下一批,且当前批次运行时会忽略新的批次触发信号(和你用全局变量的行为一致)。
核心思路
exhaustMap:接收外部Observable(这里是rx.interval(1)的批次触发信号),但仅当内部Observable(批次任务流)未在执行时,才会订阅新的内部流;若内部流正在运行,则直接忽略外部的新信号,直到内部流完成。- 批次内任务并行:用
flatMap将批次中的每个任务转换为独立的Observable,让它们同时执行。
修改后的代码
import rx from rx import operators as ops def start_task(value): print(f"Started {value}") return value def end_task(value): print(f"End: {value}") def main(): print("Start main") rx.interval(1).pipe( # 用exhaustMap替代全局变量控制批次串行 ops.exhaustMap(lambda time: rx.from_([1,2]).pipe( ops.map(lambda value: [time, value]), # 批内任务并行执行 ops.flatMap(lambda val: rx.of(val).pipe( ops.map(start_task), ops.delay(2), # 模拟任务运行时长 ops.map(end_task) )) )) ).run() if __name__ == "__main__": main()
代码说明
exhaustMap:确保只有当前批次的所有任务都完成后,才会响应下一个interval触发的批次信号,从根源避免了批次重叠。- 批内并行:通过
flatMap将每个任务包装成独立的Observable,它们会同时启动,各自延迟2秒后结束,实现批次内任务并行。 - 输出效果:和你用全局变量的方案一致,不会出现上一批未结束就启动下一批的情况。
可选:若需缓存触发信号(非忽略)
如果你的需求变为:即使批次处理时间超过1秒,也要把期间的触发信号缓存,等当前批次完成后立即启动下一批(无需等待下一个interval),可以将exhaustMap替换为**concatMap**,它会串行处理所有外部信号,按顺序执行每个批次。
内容的提问来源于stack exchange,提问作者John Ericksen
相关产品推荐
相关产品推荐

