TqSdk 异步任务怎么写?并发订阅与可维护结构

📅 2026/8/13 8:08:34
TqSdk 异步任务怎么写?并发订阅与可维护结构
TqSdk 的异步任务适合在同一个线程、同一个 API 连接中组织多个合约或多个独立监听职责。api.create_task用来注册协程api.register_update_notify为任务提供更新通知主程序仍通过wait_update推进内核。它能减少复制多份事件循环的代码但不会自动解决共享账户、重复委托和异常恢复。一个合约一个监听任务的基本结构下面为两个合约分别创建行情监听任务。每个任务独立订阅 Quote并只处理自身更新主循环统一推进数据。import os from tqsdk import TqApi, TqAuth api TqApi( authTqAuth(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_lengthwindow 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过滤自身事件。行情计算与账户交易分工清楚共享账户只有一个协调入口。保存任务引用并监控异常不吞错不让任务静默停止。退出、重启和并发冲突都有测试任务配置保留合约与参数定位。异步的价值是把多个独立等待过程组织在一起而不是让业务状态自动变安全。先把每个任务的输入、输出、账户权限和失败影响说清楚再使用协程代码才能在合约数量增加后仍然可读和可恢复。