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

Python Asyncio多队列脚本执行完成后无法终止的问题排查求助

Hey there! Let's break down why your asyncio script with multiple queues isn't terminating, and fix it step by step.

What's Wrong with Your Code?

First Version (Using queue.join())

Your consumer tasks (square_it, double_it, save_it) run in an infinite while True loop. When you call await asyncio.gather(*square_scan), the event loop gets stuck waiting for those square tasks to finish—but they never will, since they're trapped in the infinite loop. The queue.join() and task cancellation code never even gets a chance to run.

Second Version (Using Sentinel Values)

You added a single None sentinel to queue1, but you created 5 square_it tasks. Only one of those tasks will pick up the None and exit; the other 4 will keep waiting forever for more items from queue1. That's why await asyncio.gather(*square_scan) hangs—those 4 tasks are still running, waiting for input that never comes. The same problem cascades to your double_scan and save_scan tasks too.

The Fix: Proper Sentinel Handling + Correct Task Flow

We need two key adjustments:

  1. Send enough sentinels to match the number of consumer tasks for each queue (so every worker gets a signal to exit).
  2. Structure the main flow to avoid waiting for infinite tasks before processing queues.

Here's the corrected code:

import os
import warnings
import asyncio
import random
import pandas as pd
from datetime import datetime

os.environ['PYTHONASYNCIODEBUG'] = '1'
warnings.resetwarnings()

class asyncio_toy():
    def __init__(self):
        self.df = pd.DataFrame(columns=['id','final_value'])

    async def generate_random_number(self, queue1, num_consumers):
        for k in range(50):
            r = random.randint(0,10)
            await queue1.put((k, r))
        # Send one sentinel per consumer task to trigger exit
        for _ in range(num_consumers):
            await queue1.put(None)

    async def square_it(self, n, queue1, queue2, num_next_consumers):
        while True:
            print(f'{datetime.now()} - START SQUARE IT task {n} q1: {str(queue1.qsize()).zfill(2)} - q2:{str(queue2.qsize()).zfill(2)}')
            r = await queue1.get()
            if r is None:
                # Pass sentinel to the next queue's consumers
                await queue2.put(None)
                queue1.task_done()
                break
            await asyncio.sleep(5)
            await queue2.put((r[0], r[1]*r[1]))
            queue1.task_done()
            print(f'{datetime.now()} - END SQUARE IT task {n} q1: {str(queue1.qsize()).zfill(2)} - q2:{str(queue2.qsize()).zfill(2)}')

    async def double_it(self, n, queue2, queue3, num_next_consumers):
        while True:
            print(f'{datetime.now()} - START DOUBLE IT task {n} q2: {str(queue2.qsize()).zfill(2)} - q3:{str(queue3.qsize()).zfill(2)}')
            r = await queue2.get()
            if r is None:
                await queue3.put(None)
                queue2.task_done()
                break
            await asyncio.sleep(10)
            await queue3.put((r[0], 2*r[1]))
            queue2.task_done()
            print(f'{datetime.now()} - END DOUBLE IT task {n} q2: {str(queue2.qsize()).zfill(2)} - q3:{str(queue3.qsize()).zfill(2)}')

    async def save_it(self, n, queue3):
        while True:
            print(f'{datetime.now()} - START SAVE IT task {n} q3: {str(queue3.qsize()).zfill(2)}')
            r = await queue3.get()
            if r is None:
                queue3.task_done()
                break
            await asyncio.sleep(1)
            self.df.loc[len(self.df)] = [r[0], r[1]]
            self.df.to_csv('final_result.csv')
            queue3.task_done()
            print(f'{datetime.now()} - END SAVE IT task {n} q3: {str(queue3.qsize()).zfill(2)}')

    async def main(self):
        queue1 = asyncio.Queue()
        queue2 = asyncio.Queue()
        queue3 = asyncio.Queue()

        num_square_workers = 5
        num_double_workers = 5
        num_save_workers = 5

        # Create all tasks
        rand_gen = asyncio.create_task(self.generate_random_number(queue1, num_square_workers))
        square_scan = [asyncio.create_task(self.square_it(k, queue1, queue2, num_double_workers)) for k in range(num_square_workers)]
        double_scan = [asyncio.create_task(self.double_it(k, queue2, queue3, num_save_workers)) for k in range(num_double_workers)]
        save_scan = [asyncio.create_task(self.save_it(k, queue3)) for k in range(num_save_workers)]

        # Wait for producer to finish adding all data
        await rand_gen

        # Wait for all items in each queue to be fully processed
        await queue1.join()
        await queue2.join()
        await queue3.join()

        # Wait for all consumer tasks to exit after handling sentinels
        await asyncio.gather(*square_scan)
        await asyncio.gather(*double_scan)
        await asyncio.gather(*save_scan)

### testing
if __name__ == '__main__':
    toy = asyncio_toy()
    asyncio.run(toy.main())

Key Changes Explained

  1. Matching Sentinel Counts: We pass the number of consumer tasks to each stage, so we send exactly one None per worker. This ensures every task gets the exit signal.
  2. Task Flow Order: We first wait for the data producer to finish, then use queue.join() to wait for all items to be processed, then wait for consumers to exit after handling their sentinels.
  3. Proper task_done() for Sentinels: Even when a worker encounters a sentinel, we call queue.task_done() to ensure queue.join() doesn't hang.

Why Your Single-Queue Example Worked

In your single-queue code, you followed the correct order: wait for producers to finish, call queue.join() to process all items, then cancel the consumer tasks. The multi-queue version needs the same logic, but with sentinels to coordinate exit across multiple processing stages.

内容的提问来源于stack exchange,提问作者Dariva

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.29 21:22:36