如何理解并在Python中结合PostgreSQL使用asyncio?解决代码报错及输出捕获问题
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())
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}")
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

