如何用Python asyncio避免回调地狱?附服务交互场景问询
Nice scenario you've got here! Callback hell is a common pain point when dealing with asynchronous workflows, but Python's asyncio with async/await syntax is perfect for flattening this out and keeping your code clean. Let's walk through how to implement each service with asyncio, avoiding nested callbacks entirely.
The key is to replace nested callback functions with async/await syntax, which lets you write asynchronous code in a linear, synchronous-looking style. We'll use async-compatible libraries for RabbitMQ (like aio_pika) and HTTP (like aiohttp for clients, FastAPI for the Mediator's API) to keep everything running in the asyncio event loop without blocking.
1. JobInitiator: Async Task Publishing
Instead of using callback-based RabbitMQ clients, we'll use aio_pika to publish tasks asynchronously on a timer. The loop runs every X minutes with no nested callbacks—just sequential await calls.
import asyncio from aio_pika import connect, Message async def publish_task(channel, queue_name): # Simulate generating a task payload task_payload = b"sample_task_data" await channel.default_exchange.publish( Message(task_payload), routing_key=queue_name ) print(f"Published task at {asyncio.get_event_loop().time()}") async def job_initiator_loop(x_minutes, queue_name): connection = await connect("amqp://guest:guest@localhost/") async with connection: channel = await connection.channel() await channel.declare_queue(queue_name) interval = x_minutes * 60 while True: await publish_task(channel, queue_name) await asyncio.sleep(interval) if __name__ == "__main__": asyncio.run(job_initiator_loop(x_minutes=5, queue_name="task_queue"))
2. Mediator: Async REST API + RabbitMQ Consumption
The Mediator exposes an async REST endpoint using FastAPI. When an Executor asks for tasks, it asynchronously fetches from the RabbitMQ queue—no callbacks, just await for the queue get operation.
from fastapi import FastAPI import asyncio from aio_pika import connect, Queue app = FastAPI() queue: Queue = None async def setup_rabbitmq(): global queue connection = await connect("amqp://guest:guest@localhost/") channel = await connection.channel() queue = await channel.declare_queue("task_queue") @app.on_event("startup") async def startup_event(): await setup_rabbitmq() @app.get("/tasks/next") async def get_next_task(): # Asynchronously get a message from the queue (non-blocking) message = await queue.get() if message: await message.ack() return {"task": message.body.decode()} return {"task": None} if __name__ == "__main__": import uvicorn uvicorn.run(app, host="0.0.0.0", port=8000)
3. Executor: Async Polling + Task Execution
The Executor uses aiohttp to make async HTTP requests to the Mediator on a timer. Every Y minutes, it polls for a task, and if there's one, executes it asynchronously—all with linear await calls, no nested callbacks.
import asyncio import aiohttp async def fetch_task(session): async with session.get("http://localhost:8000/tasks/next") as response: return await response.json() async def execute_task(task_data): # Simulate task execution (replace with your actual logic) print(f"Executing task: {task_data}") await asyncio.sleep(5) # Simulate work print(f"Completed task: {task_data}") async def executor_loop(y_minutes): interval = y_minutes * 60 async with aiohttp.ClientSession() as session: while True: task_response = await fetch_task(session) task_data = task_response.get("task") if task_data: await execute_task(task_data) await asyncio.sleep(interval) if __name__ == "__main__": asyncio.run(executor_loop(y_minutes=2))
- Linear code flow: Instead of nesting callbacks (e.g., "on HTTP response, then on RabbitMQ message, then execute task"), we use
awaitto pause execution until each async operation completes, keeping code flat and readable. - No callback chaining: Every async operation is handled as a sequential step in the function, not as a nested function passed to another call.
- Async-compatible libraries: We use
aio_pika,aiohttp, andFastAPI—all designed to work with asyncio's event loop without blocking, so we never have to fall back to callback-based APIs.
- Use
asyncio.create_task()if you need to run multiple tasks concurrently (e.g., executing multiple tasks at once in the Executor). - Add error handling with
try/exceptblocks aroundawaitcalls to handle network issues or RabbitMQ errors gracefully. - For more complex workflows, consider using
asyncio.gather()to run multiple async operations in parallel.
内容的提问来源于stack exchange,提问作者Mr T.

