如何向asyncio任务或AnyIO TaskGroup动态添加任务?
问题:asyncio/线程/消息队列非阻塞执行优化需求
本人是asyncio、线程与消息队列的新手,在非阻塞执行方面遇到问题,现详述自研方案并寻求优化建议(这是我的首个提问,欢迎要求澄清或告知社区规则)。
我的应用用于从Broker队列消费代表待执行任务的消息,任务执行结果及生成的日志需推送至另一Broker队列。当前实现中,消息消费者运行在线程中,所有消费的消息存入共享队列;Runtime类从共享队列获取任务、构建任务参数并执行,队列空时会休眠等待新任务。
当前代码实现
启动代码
async def start(): shared_queue = Queue() futures = [] for kind, actor in ACTORS_MAP.items(): logger.info("starting") event_loop = asyncio.new_event_loop() listen_thread = Thread(target=run_loop, args=(event_loop,)) listen_thread.start() futures.append( asyncio.run_coroutine_threadsafe( ainit_and_listen(kind, shared_queue), event_loop ) ) # await asyncio.sleep(1) runtime = Runtime(shared_queue) await runtime.invoke()
监听代码
async def ainit_and_listen(actor_kind: str, runtime_queue: Queue): logger.info("Declaring queue for {0}".format(actor_kind)) connection = await connect(os.environ.get("RABBIT_URL")) channel = await connection.channel() exchange = await channel.declare_exchange(os.environ.get("CONSUMER_EXCHANGE"), "topic") queue = await channel.declare_queue(f"{os.environ.get('TEAM_ID')}.{actor_kind}") await queue.bind(exchange=exchange, routing_key=f"{os.environ.get('TEAM_ID')}.{actor_kind}") logger.info("Queue was declared and binded for {0}".format(actor_kind)) await listen(broker_queue=queue, runtime_queue=runtime_queue) await asyncio.Future() async def listen(broker_queue, runtime_queue: Queue): while True: logger.info("Listen to {0} queue".format(broker_queue.name)) incoming_message: Optional[ AbstractIncomingMessage ] = await broker_queue.get(timeout=1, fail=False) if incoming_message: # Confirm message await incoming_message.ack() data = json.loads(incoming_message.body) runtime_queue.put(data) logger.info("New message of {0} kind".format(data["actor"])) else: logger.info("No new messages. Sleeping...") await asyncio.sleep(os.environ.get("CONSUME_FREQUENCY"))
Runtime类核心方法
async def invoke(self): await self._init_pub_exchange() just_logger.info("Producer exchange was successfully declared") await self.run() async def run(self): self._running = True just_logger.info("Run starts") async with anyio.create_task_group() as tg: while True: await self.process_queue(tg) async def process_queue(self, task_group): while self._queue: data = self._queue.get() just_logger.info("Processing next message") task_group.start_soon(self.process_node, data) async def process_node(self, data): user : User = User.model_validate_json(data["user"]) workflow_id = data["workflow_id"] actor = data["actor"] arguments = data["arguments"] schema = ACTOR_SCHEMAS[actor['actor']] actor_obj = self.init_actor(actor, workflow_id) just_logger.debug("Initialized {0}".format(actor['actor'])) actor_params: dict[str, Any] = self.get_messages(schema, arguments) actor_params["context"] = self.build_context(user) actor_params["tools"] = self.build_tools(user) just_logger.debug("Set up {0}".format(actor['actor'])) actor_coro = actor_obj(**actor_params) try: just_logger.info("Executing {0}...".format(actor['actor'])) results: list[BaseMessage] | BaseMessage = await actor_coro # type: ignore except TypeError as e: just_logger.error( "{0} Actor run finished with error:\n{1}".format(actor["actor"], e) ) raise TypeError( f"Exception while running the {actor['actor']} actor: {e}" ).with_traceback(e.__traceback__) just_logger.info("Finished {0}".format(actor['actor'])) if isinstance(results, list): json_results = [result.model_dump_json() for result in results] else: json_results = [results.model_dump_json()] await self.publish_result(workflow_id, actor['node_id'], json_results) just_logger.debug("Published {0}").format(actor['actor_id'])
问题与期望
我希望在收到队列中的消息时立即启动任务,且Runtime不等待任务执行结果。我曾考虑使用AnyIO TaskGroup并传入回调参数,但任务并未执行。请问是否有办法向asyncio任务或AnyIO TaskGroup动态添加任务?
期望日志输出:
>>> New message 1 >>> New message 2 >>> New message 3 >>> New message 4 >>> Set Up 1 >>> Executing 1 >>> Set Up 2 >>> Executing 2 >>> Set Up 3 >>> Executing 3 >>> Set Up 4 >>> Executing 4 >>> Finished 1 >>> Finished 2 >>> Finished 3 >>> Finished 4
当前实际日志输出:
>>> New message 1 >>> New message 2 >>> New message 3 >>> New message 4 >>> Set Up 1 >>> Executing 1 >>> Finished 1 >>> Set Up 2 >>> Executing 2 >>> Finished 2 >>> Set Up 3 >>> Executing 3 >>> Finished 3 >>> Set Up 4 >>> Executing 4 >>> Finished 4
解决方案
问题根源
当前代码的核心问题:
- 使用同步队列(标准库
queue.Queue)而非asyncio异步队列asyncio.Queue,同步队列的get()是阻塞调用,会卡住asyncio事件循环,导致任务无法并发执行。 process_queue的轮询逻辑依赖同步队列的空判断,阻塞事件循环,导致每次只能处理一个任务,直到完成才会取下一个。
修复步骤
1. 替换为异步队列
将启动代码中的同步队列替换为asyncio.Queue,并修改所有队列操作为异步调用:
- 启动代码修改:
async def start(): shared_queue = asyncio.Queue() # 改用asyncio异步队列 futures = [] for kind, actor in ACTORS_MAP.items(): logger.info("starting") event_loop = asyncio.new_event_loop() listen_thread = Thread(target=run_loop, args=(event_loop,)) listen_thread.start() futures.append( asyncio.run_coroutine_threadsafe( ainit_and_listen(kind, shared_queue), event_loop ) ) # await asyncio.sleep(1) runtime = Runtime(shared_queue) await runtime.invoke()
- 监听代码修改:
async def listen(broker_queue, runtime_queue: asyncio.Queue): while True: logger.info("Listen to {0} queue".format(broker_queue.name)) incoming_message: Optional[ AbstractIncomingMessage ] = await broker_queue.get(timeout=1, fail=False) if incoming_message: await incoming_message.ack() data = json.loads(incoming_message.body) await runtime_queue.put(data) # 异步put操作 logger.info("New message of {0} kind".format(data["actor"])) else: logger.info("No new messages. Sleeping...") await asyncio.sleep(int(os.environ.get("CONSUME_FREQUENCY", 1))) # 字符串转整数
2. 修改Runtime队列处理逻辑
删除原process_queue方法,改为异步等待队列消息,让事件循环能正常调度并发任务:
async def run(self): self._running = True just_logger.info("Run starts") async with anyio.create_task_group() as tg: while self._running: # 异步等待队列消息,不阻塞事件循环 data = await self._queue.get() just_logger.info("Processing next message") tg.start_soon(self.process_node, data) # 标记任务已处理(可选,用于队列join操作) self._queue.task_done()
3. 验证TaskGroup动态添加能力
AnyIO的TaskGroup.start_soon()本身支持动态添加任务,之前未生效是因为事件循环被同步队列阻塞。换成异步队列后,事件循环可正常调度,新任务会被立即启动,无需等待前一个任务完成。
其他优化建议
- 移除不必要的线程:如果Broker客户端支持asyncio(如aio-pika),可直接在主事件循环中运行消费者协程,减少线程切换开销。
- 增强错误处理:在
process_node中捕获通用异常,避免单个任务失败导致整个TaskGroup崩溃(AnyIO TaskGroup默认会因单个任务异常取消所有任务)。 - 配置参数校验:对环境变量做类型转换和默认值设置,避免类型错误。
内容的提问来源于stack exchange,提问作者Mini
相关产品推荐
相关产品推荐

