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

如何用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.

Core Approach to Avoid Callback Hell

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))
Why This Avoids Callback Hell
  • Linear code flow: Instead of nesting callbacks (e.g., "on HTTP response, then on RabbitMQ message, then execute task"), we use await to 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, and FastAPI—all designed to work with asyncio's event loop without blocking, so we never have to fall back to callback-based APIs.
Bonus Tips
  • 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/except blocks around await calls 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.

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:05:13