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

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()

代码说明

  1. exhaustMap:确保只有当前批次的所有任务都完成后,才会响应下一个interval触发的批次信号,从根源避免了批次重叠。
  2. 批内并行:通过flatMap将每个任务包装成独立的Observable,它们会同时启动,各自延迟2秒后结束,实现批次内任务并行。
  3. 输出效果:和你用全局变量的方案一致,不会出现上一批未结束就启动下一批的情况。

可选:若需缓存触发信号(非忽略)

如果你的需求变为:即使批次处理时间超过1秒,也要把期间的触发信号缓存,等当前批次完成后立即启动下一批(无需等待下一个interval),可以将exhaustMap替换为**concatMap**,它会串行处理所有外部信号,按顺序执行每个批次。

内容的提问来源于stack exchange,提问作者John Ericksen

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 17:25:01