You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何向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

解决方案

问题根源

当前代码的核心问题:

  1. 使用同步队列(标准库queue.Queue)而非asyncio异步队列asyncio.Queue,同步队列的get()是阻塞调用,会卡住asyncio事件循环,导致任务无法并发执行。
  2. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.26 05:57:03