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:
- Send enough sentinels to match the number of consumer tasks for each queue (so every worker gets a signal to exit).
- 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
- Matching Sentinel Counts: We pass the number of consumer tasks to each stage, so we send exactly one
Noneper worker. This ensures every task gets the exit signal. - 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. - Proper
task_done()for Sentinels: Even when a worker encounters a sentinel, we callqueue.task_done()to ensurequeue.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

