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
相关产品推荐
相关产品推荐

