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

RxPy:如何从外部回调创建热可观察对象并订阅多异步进程?

使用RxPy将外部服务回调转换为Observable并支持多异步订阅

最近我碰到了这么个需求:要对接一个叫ExternalDummyService的外部服务,这个服务得注册回调函数才能接收事件,我想把这些回调事件转换成RxPy的Observable,方便多个异步进程来订阅处理。下面是整理好的核心实现思路和代码:

核心实现代码

from rx import Subject
import threading

class ExternalDummyService:
    def __init__(self):
        self.callback = None
        self.subject = Subject()
    
    def register_callback(self, callback):
        self.callback = callback
    
    def start(self):
        def run():
            # 用循环模拟耗时操作(替代sleep,因为测试环境里sleep可能没法正常工作)
            foo = 0
            for i in range(10000):
                foo += i
            # 触发回调推送数据
            if self.callback:
                self.callback("来自外部服务的事件数据")
        threading.Thread(target=run).start()

# 初始化外部服务实例
thread = ExternalDummyService()

# 创建可连接Observable,用publish()让多个订阅者共享同一份数据流
external_obs = thread.subject.publish()

# 绑定回调:把外部服务的回调事件推送到Subject中
thread.register_callback(lambda data: thread.subject.on_next(data))

# 第一个异步订阅者
external_obs.subscribe(lambda x: print(f"订阅者1收到数据: {x}"))

# 第二个异步订阅者
external_obs.subscribe(lambda x: print(f"订阅者2收到数据: {x}"))

# 激活可连接的Observable,让所有订阅者开始接收数据
external_obs.connect()

# 启动外部服务
thread.start()

关键细节说明

  • Subject的桥梁作用:Subject在这里既是Observable也是Observer,把外部服务的回调和它的on_next方法绑定后,外部服务触发回调时,就能自动把数据推送给所有订阅者。
  • publish()的必要性:默认的Observable是"冷"的,每个订阅者都会单独触发一次数据流。而我们需要"热"Observable——让所有订阅者共享同一份外部服务事件,所以用publish()创建可连接Observable,最后通过connect()激活,确保所有订阅者收到的是同一批事件。
  • 耗时操作模拟:因为测试环境里sleep可能无法正常工作,所以用循环累加的方式模拟异步任务的耗时,保证外部服务的回调是在异步线程中触发的。

内容的提问来源于stack exchange,提问作者raul.vila

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 10:05:17