Python如何将拉取Bloomberg实时数据的async for循环封装为线程可用函数
Bloomberg实时订阅异步逻辑封装为线程运行函数实现
原有代码问题修正
你贴的原代码有两个会直接导致运行失败的错误:
- 用Python内置类型名
list作为自定义变量名,会覆盖内置列表方法,引发不可预期的错误 - 列表的
append()方法是原地修改操作,返回值为None,原写法list= list.append(d['LAST_PRICE'])执行一次后就会把变量置为None,后续追加数据会直接抛异常
封装思路
异步代码无法直接在普通线程函数中运行,需要在线程入口内独立创建asyncio事件循环,把异步订阅逻辑放到事件循环中驱动执行;跨线程读写存储的价格数据时加线程锁避免竞态问题。
可直接复用的实现代码
import threading import asyncio # 存储实时价格的容器 last_price_records = [] # 跨线程读写数据用的锁 data_lock = threading.Lock() def blp_live_sub_thread(blp_instance, security='IBM US Equity', fields=None): """ 可放入线程运行的Bloomberg实时数据订阅函数 :param blp_instance: 已经初始化好的blp连接实例 :param security: 订阅的证券代码 :param fields: 订阅的字段列表 """ if fields is None: fields = ['LAST_PRICE'] # 为当前线程创建独立的事件循环 loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) async def _async_sub(): async for data in blp_instance.live(security, flds=fields, info=fields): price = data['LAST_PRICE'] # 加锁写入数据,避免多线程读写冲突 with data_lock: last_price_records.append(price) # 此处可扩展自定义逻辑,比如触发价格预警、推送到消息队列等 try: # 阻塞运行异步订阅,直到订阅被取消/连接断开 loop.run_until_complete(_async_sub()) finally: # 线程退出前清理事件循环资源 loop.close() # 启动线程示例 if __name__ == '__main__': # 提前初始化blp连接,注意blp实例不要跨线程传递调用 # blp = 你的Bloomberg连接初始化逻辑 # 标记为守护线程,主程序退出时自动终止订阅线程 sub_thread = threading.Thread( target=blp_live_sub_thread, args=(blp,), daemon=True ) sub_thread.start() # 其他业务逻辑示例:每秒打印最新价格 import time while True: time.sleep(1) with data_lock: if last_price_records: print(f"最新成交价: {last_price_records[-1]}")
关键注意事项
- Bloomberg的连接实例
blp不要跨线程调用,所有和订阅相关的操作都要在工作线程的事件循环内执行,否则会触发asyncio线程安全错误 - 如果需要主动停止订阅,不要强制杀死线程,通过
loop.call_soon_threadsafe()方法向事件循环发送协程取消信号即可安全退出 - 如果需要在多个地方消费价格数据,可以把列表换成线程安全的队列,避免手动加锁的复杂度
内容的提问来源于stack exchange,提问作者Gauri Kant
相关产品推荐
相关产品推荐

