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

如何理解并在Python中结合PostgreSQL使用asyncio?解决代码报错及输出捕获问题

Fixing "event loop stopped before the future completed" Error

This error almost always means your event loop is exiting before all async tasks/futures have finished running. Let’s break down the most common scenarios and how to fix them:

Scenario 1: Manual loop management gone wrong

If you’re handling the loop yourself (instead of using asyncio.run()), calling loop.stop() too early will cut off pending tasks. Instead, wait for all tasks to complete before shutting down the loop.

Bad example:

import asyncio

async def slow_task():
    await asyncio.sleep(2)
    print("Task finished")

loop = asyncio.get_event_loop()
loop.create_task(slow_task())
loop.stop()  # Loop stops before the task finishes
loop.run_forever()

Fixed version: Use run_until_complete to wait for your task, or let asyncio.run() handle lifecycle management:

import asyncio

async def slow_task():
    await asyncio.sleep(2)
    print("Task finished")

# Modern, recommended approach
asyncio.run(slow_task())

Scenario 2: Forgetting to await a coroutine

If you create a task but never await it, the loop might exit before it completes. Always ensure you await every coroutine or task you care about.

Bad example:

async def main():
    # Creates a task but doesn't wait for it
    asyncio.create_task(slow_task())

asyncio.run(main())  # Loop exits immediately

Fixed version: Await the task directly, or use asyncio.gather for multiple tasks:

async def main():
    task = asyncio.create_task(slow_task())
    await task  # Wait for the task to finish

asyncio.run(main())

Capturing Output from Asyncio Code

Async functions return Future objects (or coroutines that wrap into futures), so capturing output is as simple as awaiting them. Here’s how to do it in different scenarios:

Single task output

Assign the awaited result to a variable directly:

async def calculate_total(a, b):
    await asyncio.sleep(1)
    return a + b

async def main():
    total = await calculate_total(10, 20)
    print(f"Total: {total}")  # Outputs "Total: 30"

asyncio.run(main())

Multiple concurrent tasks

Use asyncio.gather() to run tasks in parallel and collect all results in a list:

async def fetch_resource(resource_id):
    await asyncio.sleep(1)
    return f"Resource {resource_id} data"

async def main():
    tasks = [fetch_resource(1), fetch_resource(2), fetch_resource(3)]
    results = await asyncio.gather(*tasks)
    for res in results:
        print(res)

asyncio.run(main())

Output from long-running background tasks

If you need to capture results as tasks finish (instead of waiting for all), use asyncio.as_completed():

async def main():
    tasks = [fetch_resource(i) for i in range(3)]
    for future in asyncio.as_completed(tasks):
        result = await future
        print(f"Received: {result}")

Using Asyncio with PostgreSQL

To work with PostgreSQL asynchronously, you’ll need an async driver—asyncpg is the industry standard. Let’s cover the core workflows:

Step 1: Install asyncpg

First, install the package via pip:

pip install asyncpg

Step 2: Basic Connection & Query

Use async context managers to handle connection setup/cleanup automatically:

import asyncio
import asyncpg

async def main():
    # Connect to your database
    conn = await asyncpg.connect(
        user="your_username",
        password="your_password",
        database="your_db",
        host="localhost"
    )

    try:
        # Fetch multiple rows
        users = await conn.fetch("SELECT id, username FROM users LIMIT 5")
        for user in users:
            print(f"User: {user['username']} (ID: {user['id']})")
        
        # Fetch a single row
        admin = await conn.fetchrow("SELECT * FROM users WHERE is_admin = $1", True)
        if admin:
            print(f"Admin found: {admin['username']}")
        
        # Execute a write operation
        await conn.execute("INSERT INTO users (username) VALUES ($1)", "new_user_123")
    finally:
        # Close the connection
        await conn.close()

asyncio.run(main())

Step 3: Connection Pools (Production Best Practice)

For production, use a connection pool to reuse connections efficiently:

import asyncio
import asyncpg

async def main():
    # Create a pool with min/max connections
    pool = await asyncpg.create_pool(
        user="your_username",
        password="your_password",
        database="your_db",
        host="localhost",
        min_size=5,
        max_size=15
    )

    # Acquire a connection from the pool
    async with pool.acquire() as conn:
        count = await conn.fetchval("SELECT COUNT(*) FROM users")
        print(f"Total users: {count}")
    
    # Close the pool when done
    await pool.close()

asyncio.run(main())

Step 4: Async Transactions

Use async transactions to ensure atomicity for multi-step operations:

async def transfer_money(conn, from_account, to_account, amount):
    async with conn.transaction():
        # Deduct from sender
        await conn.execute(
            "UPDATE accounts SET balance = balance - $1 WHERE id = $2",
            amount, from_account
        )
        # Add to receiver
        await conn.execute(
            "UPDATE accounts SET balance = balance + $1 WHERE id = $2",
            amount, to_account
        )

This ensures both queries either succeed together or roll back entirely if an error occurs.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 04:19:28