TqSdk 的异步任务适合在同一个线程、同一个 API 连接中组织多个合约或多个独立监听职责。api.create_task用来注册协程,api.register_update_notify为任务提供更新通知,主程序仍通过wait_update推进内核。它能减少复制多份事件循环的代码,但不会自动解决共享账户、重复委托和异常恢复。
一个合约一个监听任务的基本结构
下面为两个合约分别创建行情监听任务。每个任务独立订阅 Quote,并只处理自身更新;主循环统一推进数据。
import os from tqsdk import TqApi, TqAuth api = TqApi( auth=TqAuth(os.environ["TQ_USER"], os.environ["TQ_PASSWORD"]) ) async def watch_quote(symbol): quote = await api.get_quote(symbol) async with api.register_update_notify() as updates: async for _ in updates: if api.is_changing(quote, "last_price"): print(symbol, quote.datetime, quote.last_price) api.create_task(watch_quote("SHFE.rb2610")) api.create_task(watch_quote("DCE.m2609")) try: while True: api.wait_update() finally: api.close()异步调用await api.get_quote会等待合约行情准备后返回,更适合协程内部使用。更新通道可能收到许多与当前对象无关的通知,因此任务仍要用is_changing过滤真正关心的字段。
create_task 只负责调度,不负责业务隔离
多个任务在一个线程中协作运行,某个任务只有在await处让出执行机会。若在协程中做长时间同步计算、文件写入或阻塞网络请求,其他任务和主事件循环都会受到影响。耗时工作应缩小、分批,或交给清楚的外部执行边界。
任务之间共享同一个api,也可能共享同一账户。两个任务若操作同一合约,各自的信号都可能创建订单。把代码拆成两个函数并没有形成交易隔离,账户动作仍需要统一协调。
可以让行情任务只产出信号,让一个交易管理任务读取信号并检查账户、持仓与活动委托。这样每个合约的计算可以独立,最终订单仍有唯一入口。
任务之间传递的信息应尽量小,例如合约、信号值、生成时间和规则版本,而不是直接共享一大块可变字典。消费端收到后先检查时间是否过期,再与当前账户状态组合。这样任务处理有延迟时,不会把旧信号当成新指令。
共享变量若不可避免,应规定唯一写入者。多个协程同时修改目标仓位、最后处理时间或活动委托表,虽然在单线程中不会像多线程那样同时执行一条指令,仍可能在await前后形成难以预期的先后关系。
更新通知怎样避免无效工作
register_update_notify提供的是更新发生通知,不保证每次都与本任务对象有关。任务内部先检查is_changing,再做计算。K 线策略若只在新线形成时运行,就监听末行datetime,不要每次通知都重算整段指标。
async def watch_kline(symbol, window): klines = await api.get_kline_serial( symbol, 60, data_length=window + 2 ) async with api.register_update_notify() as updates: async for _ in updates: if api.is_changing(klines.iloc[-1], "datetime"): closed = klines.iloc[:-1].copy() signal = calculate_signal(closed, window) publish_signal(symbol, signal)两个普通函数是结构占位。它们不应在内部直接创建另一套 API,也不应在同一信号未处理完时无限堆积新任务。若信号生产速度超过消费速度,需要设定覆盖、排队或丢弃旧信号的明确策略。
异常不能悄悄留在后台
后台任务发生异常后,若主循环从不观察任务结果,程序可能继续运行,却已经失去一个合约的处理能力。创建任务时应保存任务引用,建立统一的异常监控与日志,让任何任务失败都能定位到任务名和合约。
异常策略取决于任务依赖。完全独立的行情展示任务失败,其他任务可能继续;参与同一组合策略的任一任务失败,通常应暂停整组新动作。依赖关系应由配置明确,而不是在异常发生后临时判断。
不要在协程里用宽泛except Exception: pass。它会吞掉字段错误、业务错误和连接问题,任务表面存活却不再产生有效结果。可恢复异常应记录并重试,不可解释异常应停止相关动作并向主控报告。
任务引用可以附带清楚名称与合约,监控层定期检查是否结束。任务正常结束也要区分是设计完成还是意外提前返回;长期监听任务无故结束,和抛出异常一样需要处理。
重试必须有限并带有状态复核。协程失败后直接创建一份新任务,旧任务可能仍拥有活动委托或未消费消息。重新启动前先清理旧任务影响,确认账户事实,再恢复订阅。
任务数量与连接数量怎样选择
少量独立策略最简单的方式可能仍是一进程一个实例,启动、停止和调试都直观;但每个进程会建立独立连接,数量增加后连接与资源成本上升。单线程多异步任务只使用一份连接,适合结构相似、能够共享事件循环的实例。
选择依据不是“异步一定更快”,而是任务是否能短时间完成每次处理、是否共享账户、是否需要独立故障边界。复杂策略若团队无法稳定调试协程,简单进程隔离可能更可靠。
无论哪种方式,都要记录合约、参数、账户和任务标识。复制多个脚本但不记录配置,与创建许多协程但没有任务名,都会让运行现场难以追踪。
异步并不适合所有计算。大量 pandas 运算、模型推理或同步磁盘操作会占用事件循环时间,拆成协程也不会自动并行。先测量单次处理耗时;超过行情容忍范围时,使用独立进程或工作队列,并为返回结果设置超时和版本校验。
任务数量也不应只由合约数量决定。一个任务可以负责一个合约,也可以负责一类轻量职责。选择时看失败是否需要独立、状态是否共享、日志是否可定位,而不是机械地“一合约一任务”。
退出和重启要管理所有任务
主程序退出时先停止产生新业务动作,再关闭 API,让关联任务结束。若任务负责订单,还应按既定规则处理活动委托和持仓。强制结束主进程不会自动完成账户收尾。
重启后任务内存状态全部丢失。信号缓存、已处理时间和目标仓位需要从持久化记录及账户事实恢复。尤其不能因为新任务刚创建,就假设账户里没有旧订单。
测试时主动让一个任务抛出异常,确认主控能发现;让一个任务执行较慢,确认其他行情是否延迟;让两个任务同时产生信号,确认交易入口能防止冲突。只有顺利路径的并发测试无法证明结构可维护。
异步任务清单
- 一个 API 与主
wait_update循环推进全部任务,不在协程内创建隐蔽连接。 - 每个任务通过更新通道等待,并用
is_changing过滤自身事件。 - 行情计算与账户交易分工清楚,共享账户只有一个协调入口。
- 保存任务引用并监控异常,不吞错,不让任务静默停止。
- 退出、重启和并发冲突都有测试,任务配置保留合约与参数定位。
异步的价值是把多个独立等待过程组织在一起,而不是让业务状态自动变安全。先把每个任务的输入、输出、账户权限和失败影响说清楚,再使用协程,代码才能在合约数量增加后仍然可读和可恢复。