从非异步函数后台运行asyncio协程的问题及解决方案问询
异步回调任务无法启动的问题与解决
问题描述
以下是简化后的代码:
def config_changed(new_config: Any): apply_config(new_config) start() asyncio.create_task(upload_to_db({"running": True})) async def upload_to_db(data: Any): await some_db_code(data) def start(): asyncio.create_task(run()) async def run(): while True: do_something_every_second() await asyncio.sleep(1) loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) some_db_client.listen_to_change(path_to_config, callback=config_changed) loop.run_forever()
运行时在asyncio.create_task(run())行报错:no running loop。尝试修改所有asyncio.create_task为以下代码后,任务始终无法启动(已确认start函数被调用):
try: loop = asyncio.get_event_loop() except RuntimeError: loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) loop.create_task(...)
再修改为以下代码后,任务可以运行,但upload_to_db从未被调用,因为run_until_complete会阻塞后续代码:
task = loop.create_task(...) loop.run_until_complete(task)
补充说明:start必须为非异步函数,因为它是抽象类的抽象方法,子类既可以用它运行异步函数,也可以用来注册非异步回调(如监听数据库变更),必须兼容两种场景。
解决方案
核心问题在于:config_changed作为回调函数,可能在非事件循环线程被调用,此时直接调用asyncio.create_task或重新创建loop都会导致任务无法正确绑定到已运行的事件循环上。正确的实现方式如下:
1. 提前保存全局事件循环实例
初始化时就保存已经创建并准备启动的loop,避免在回调中重新获取或创建新loop:
# 全局保存已初始化的事件循环 global_loop = asyncio.new_event_loop() asyncio.set_event_loop(global_loop)
2. 在回调中安全提交异步任务
如果回调在非事件循环线程执行,必须使用loop.call_soon_threadsafe来提交任务(create_task不是线程安全的,不能跨线程调用)。对于异步函数,需要用create_task包装后通过call_soon_threadsafe提交:
def config_changed(new_config: Any): apply_config(new_config) start() # 线程安全地提交upload_to_db任务 global_loop.call_soon_threadsafe( lambda: global_loop.create_task(upload_to_db({"running": True})) ) def start(): # 线程安全地提交run任务 global_loop.call_soon_threadsafe( lambda: global_loop.create_task(run()) )
完整修正后的代码
import asyncio from typing import Any def apply_config(config: Any): # 实现配置应用逻辑 pass async def some_db_code(data: Any): # 实现数据库操作逻辑 pass def do_something_every_second(): # 实现每秒执行的逻辑 pass def config_changed(new_config: Any): apply_config(new_config) start() # 线程安全提交异步任务 global_loop.call_soon_threadsafe( lambda: global_loop.create_task(upload_to_db({"running": True})) ) async def upload_to_db(data: Any): await some_db_code(data) def start(): # 线程安全提交异步任务 global_loop.call_soon_threadsafe( lambda: global_loop.create_task(run()) ) async def run(): while True: do_something_every_second() await asyncio.sleep(1) # 全局初始化事件循环 global_loop = asyncio.new_event_loop() asyncio.set_event_loop(global_loop) # 假设some_db_client是第三方库的客户端 some_db_client.listen_to_change(path_to_config, callback=config_changed) # 启动事件循环 global_loop.run_forever()
关键说明
call_soon_threadsafe是asyncio提供的线程安全接口,用于在事件循环线程中执行指定函数,确保异步任务能被正确添加到已运行的loop中。- 避免在回调中创建新的事件循环:新创建的loop未启动,因此提交的任务永远不会执行;而
run_until_complete会阻塞当前线程,导致后续代码无法执行。 - 保持
start的非异步特性:通过线程安全的方式提交异步任务,完美兼容子类的两种实现场景。
内容的提问来源于stack exchange,提问作者Jeffrey Chen
相关产品推荐
相关产品推荐

