Python asyncio
Asynchronous I/O framework for writing concurrent code using async/await syntax.
Python asyncio
Asynchronous I/O framework for writing concurrent code using async/await syntax.
Overview
asyncio is Python's standard library for writing concurrent code using coroutines, multiplexing I/O access, and running network clients and servers. It provides an event loop that manages asynchronous operations, enabling efficient handling of I/O-bound tasks without threading overhead.
Use asyncio for:
- Network servers and clients (HTTP, WebSocket, TCP/UDP)
- Database operations with async drivers
- Concurrent API calls and web scraping
- File I/O operations (with aiofiles)
- Long-running tasks that involve waiting
graph TB
subgraph Event Loop
A[Event Loop] --> B[Task Queue]
B --> C[Ready Tasks]
B --> D[Waiting Tasks]
C --> E[Execute Coroutine]
E --> F{I/O Operation?}
F -->|Yes| D
F -->|No| G[Continue Execution]
G --> H{Complete?}
H -->|No| C
H -->|Yes| I[Return Result]
D --> J[I/O Ready]
J --> C
end
Event Loop Fundamentals
The event loop is the core of asyncio, managing and distributing the execution of asynchronous tasks.
Running the Event Loop
import asyncio
# Recommended approach
asyncio.run(main())
# Lower-level control (get_event_loop() is deprecated when no loop is running)
loop = asyncio.new_event_loop()
try:
loop.run_until_complete(main())
finally:
loop.close()
# Get running loop (inside async context)
loop = asyncio.get_running_loop()
Event Loop Operations
import asyncio
async def main():
# Schedule callback
loop = asyncio.get_running_loop()
loop.call_soon(callback, *args)
# Schedule callback after delay
loop.call_later(5.0, callback, *args)
# Schedule at specific time
loop.call_at(loop.time() + 5.0, callback)
# Run in executor (for blocking operations)
result = await loop.run_in_executor(None, blocking_function, arg)
# Run in process pool (CPU-bound work; threads won't help under the GIL)
from concurrent.futures import ProcessPoolExecutor
executor = ProcessPoolExecutor(max_workers=3)
result = await loop.run_in_executor(executor, cpu_bound_task)
# Callback example
def callback(arg):
print(f"Callback executed with {arg}")
Event Loop Policies
Deprecated: the entire event loop policy system is deprecated since Python 3.14 and slated for removal in Python 3.16. Prefer
asyncio.run()orasyncio.Runnerwith aloop_factoryto select a loop implementation. The code below remains valid on Python 3.11–3.13.
import asyncio
# Get current policy
policy = asyncio.get_event_loop_policy()
# Set custom policy
asyncio.set_event_loop_policy(asyncio.WindowsProactorEventLoopPolicy())
# Create new event loop
new_loop = asyncio.new_event_loop()
asyncio.set_event_loop(new_loop)
Creating and Awaiting Coroutines
Coroutines are the building blocks of asyncio, defined with async def and executed with await.
stateDiagram-v2
[*] --> Created: async def declared
Created --> Scheduled: await/create_task
Scheduled --> Running: Event loop executes
Running --> Suspended: await I/O
Suspended --> Running: I/O complete
Running --> Completed: Return value
Running --> Failed: Exception raised
Completed --> [*]
Failed --> [*]
Basic Coroutines
import asyncio
# Define a coroutine
async def fetch_data(url):
print(f"Fetching {url}")
await asyncio.sleep(1) # Simulate I/O
return f"Data from {url}"
# Await a coroutine
async def main():
result = await fetch_data("https://api.example.com")
print(result)
asyncio.run(main())
Multiple Coroutines
import asyncio
async def task1():
await asyncio.sleep(1)
return "Task 1 complete"
async def task2():
await asyncio.sleep(2)
return "Task 2 complete"
async def main():
# Sequential execution (3 seconds total)
result1 = await task1()
result2 = await task2()
# Concurrent execution (2 seconds total)
results = await asyncio.gather(task1(), task2())
print(results) # ['Task 1 complete', 'Task 2 complete']
asyncio.run(main())
Coroutine Utilities
import asyncio
# Check if object is a coroutine
if asyncio.iscoroutine(obj):
result = await obj
# Check if function is a coroutine function
# (asyncio.iscoroutinefunction is deprecated since Python 3.14,
# removal planned for 3.16 -- use inspect.iscoroutinefunction instead)
import inspect
if inspect.iscoroutinefunction(func):
result = await func()
# Note: generator-based coroutines (@asyncio.coroutine) were removed in
# Python 3.11; always use async def
Task Management and Cancellation
Tasks wrap coroutines and allow them to run concurrently in the event loop.
Creating Tasks
import asyncio
async def background_task(name):
print(f"{name} started")
await asyncio.sleep(2)
print(f"{name} completed")
return f"Result from {name}"
async def main():
# Create task - starts running immediately
task1 = asyncio.create_task(background_task("Task 1"))
task2 = asyncio.create_task(background_task("Task 2"))
# Alternative: ensure_future (works with coroutines and futures)
task3 = asyncio.ensure_future(background_task("Task 3"))
# Wait for tasks
result1 = await task1
result2 = await task2
result3 = await task3
print(result1, result2, result3)
asyncio.run(main())
Task Cancellation
import asyncio
async def cancellable_task():
try:
print("Task starting")
await asyncio.sleep(10)
print("Task completed")
except asyncio.CancelledError:
print("Task was cancelled")
# Cleanup code here
raise # Re-raise to mark task as cancelled
async def main():
task = asyncio.create_task(cancellable_task())
await asyncio.sleep(1)
# Cancel the task
task.cancel()
try:
await task
except asyncio.CancelledError:
print("Confirmed task cancellation")
# Check if cancelled
print(f"Cancelled: {task.cancelled()}")
asyncio.run(main())
Task Management Patterns
import asyncio
async def worker(name, queue):
while True:
item = await queue.get()
if item is None: # Sentinel value
break
print(f"{name} processing {item}")
await asyncio.sleep(1)
queue.task_done()
async def main():
queue = asyncio.Queue()
# Create worker tasks
tasks = []
for i in range(3):
task = asyncio.create_task(worker(f"Worker-{i}", queue))
tasks.append(task)
# Add work items
for i in range(10):
await queue.put(f"Item-{i}")
# Wait for all items to be processed
await queue.join()
# Cancel workers
for task in tasks:
task.cancel()
# Wait for cancellation
await asyncio.gather(*tasks, return_exceptions=True)
asyncio.run(main())
Waiting for Multiple Tasks
import asyncio
async def task(name, delay):
await asyncio.sleep(delay)
return f"{name} done"
async def main():
tasks = [
asyncio.create_task(task("A", 1)),
asyncio.create_task(task("B", 2)),
asyncio.create_task(task("C", 3))
]
# Wait for all tasks
results = await asyncio.gather(*tasks)
print(results)
# Wait for all, return exceptions instead of raising
results = await asyncio.gather(*tasks, return_exceptions=True)
# Wait for first completion
done, pending = await asyncio.wait(
tasks,
return_when=asyncio.FIRST_COMPLETED
)
# Cancel remaining tasks
for t in pending:
t.cancel()
# Wait for first exception or all completed
done, pending = await asyncio.wait(
tasks,
return_when=asyncio.FIRST_EXCEPTION
)
# Wait with timeout
try:
results = await asyncio.wait_for(
asyncio.gather(*tasks),
timeout=2.0
)
except asyncio.TimeoutError:
print("Tasks timed out")
asyncio.run(main())
Task Groups (Python 3.11+)
import asyncio
async def task(name):
await asyncio.sleep(1)
return f"{name} completed"
async def main():
# Automatically manages task lifecycle
async with asyncio.TaskGroup() as group:
task1 = group.create_task(task("Task 1"))
task2 = group.create_task(task("Task 2"))
task3 = group.create_task(task("Task 3"))
# All tasks completed here
print(task1.result(), task2.result(), task3.result())
# If any task raises exception, all are cancelled
try:
async with asyncio.TaskGroup() as group:
group.create_task(failing_task())
group.create_task(task("Task"))
except ExceptionGroup as eg:
print(f"Task group failed: {eg}")
asyncio.run(main())
Async I/O Operations
Network I/O
import asyncio
# TCP Server
async def handle_client(reader, writer):
data = await reader.read(1024)
message = data.decode()
addr = writer.get_extra_info('peername')
print(f"Received {message} from {addr}")
response = f"Echo: {message}"
writer.write(response.encode())
await writer.drain()
writer.close()
await writer.wait_closed()
async def start_server():
server = await asyncio.start_server(
handle_client,
'127.0.0.1',
8888
)
addr = server.sockets[0].getsockname()
print(f"Serving on {addr}")
async with server:
await server.serve_forever()
# TCP Client
async def tcp_client():
reader, writer = await asyncio.open_connection(
'127.0.0.1',
8888
)
message = "Hello, server!"
writer.write(message.encode())
await writer.drain()
data = await reader.read(1024)
print(f"Received: {data.decode()}")
writer.close()
await writer.wait_closed()
asyncio.run(tcp_client())
UDP Sockets
import asyncio
class UDPProtocol(asyncio.DatagramProtocol):
def __init__(self, on_message):
self.on_message = on_message
def connection_made(self, transport):
self.transport = transport
def datagram_received(self, data, addr):
self.on_message(data, addr)
def error_received(self, exc):
print(f"Error: {exc}")
async def udp_server():
def on_message(data, addr):
print(f"Received {data.decode()} from {addr}")
loop = asyncio.get_running_loop()
transport, protocol = await loop.create_datagram_endpoint(
lambda: UDPProtocol(on_message),
local_addr=('127.0.0.1', 9999)
)
try:
await asyncio.sleep(3600) # Run for 1 hour
finally:
transport.close()
async def udp_client():
loop = asyncio.get_running_loop()
transport, protocol = await loop.create_datagram_endpoint(
lambda: UDPProtocol(lambda d, a: None),
remote_addr=('127.0.0.1', 9999)
)
transport.sendto(b"Hello, UDP server!")
await asyncio.sleep(1)
transport.close()
asyncio.run(udp_client())
HTTP Requests (with aiohttp)
import asyncio
import aiohttp
async def fetch_url(session, url):
async with session.get(url) as response:
return await response.text()
async def fetch_multiple():
async with aiohttp.ClientSession() as session:
urls = [
'https://api.github.com/users/octocat',
'https://api.github.com/users/torvalds',
'https://api.github.com/users/gvanrossum'
]
# Fetch concurrently
tasks = [fetch_url(session, url) for url in urls]
results = await asyncio.gather(*tasks)
return results
# POST request
async def post_data():
async with aiohttp.ClientSession() as session:
data = {'key': 'value'}
async with session.post('https://httpbin.org/post', json=data) as resp:
return await resp.json()
asyncio.run(fetch_multiple())
File I/O (with aiofiles)
import asyncio
import aiofiles
async def read_file(filename):
async with aiofiles.open(filename, mode='r') as f:
contents = await f.read()
return contents
async def write_file(filename, data):
async with aiofiles.open(filename, mode='w') as f:
await f.write(data)
async def process_files():
# Read multiple files concurrently
files = ['file1.txt', 'file2.txt', 'file3.txt']
tasks = [read_file(f) for f in files]
contents = await asyncio.gather(*tasks)
# Write results
await write_file('output.txt', '\n'.join(contents))
# Line-by-line reading
async def read_lines(filename):
async with aiofiles.open(filename, mode='r') as f:
async for line in f:
print(line.strip())
asyncio.run(process_files())
Database Operations (with asyncpg)
import asyncio
import asyncpg
async def fetch_users():
# Create connection pool
pool = await asyncpg.create_pool(
user='user',
password='password',
database='mydb',
host='127.0.0.1'
)
# Execute query
async with pool.acquire() as conn:
rows = await conn.fetch('SELECT * FROM users WHERE active = $1', True)
for row in rows:
print(row['name'], row['email'])
# Execute multiple queries concurrently
async with pool.acquire() as conn:
async with conn.transaction():
await conn.execute(
'INSERT INTO users(name, email) VALUES($1, $2)',
'John Doe', 'john@example.com'
)
await conn.execute(
'UPDATE accounts SET balance = balance - $1 WHERE user_id = $2',
100, 42
)
await pool.close()
asyncio.run(fetch_users())
Synchronisation Primitives
Lock
import asyncio
async def worker(lock, resource, worker_id):
async with lock:
print(f"Worker {worker_id} acquired lock")
resource['value'] += 1
await asyncio.sleep(1)
print(f"Worker {worker_id} releasing lock")
async def main():
lock = asyncio.Lock()
resource = {'value': 0}
await asyncio.gather(*[
worker(lock, resource, i) for i in range(3)
])
print(f"Final value: {resource['value']}")
asyncio.run(main())
Semaphore
import asyncio
async def download(sem, url):
async with sem:
print(f"Downloading {url}")
await asyncio.sleep(2)
print(f"Finished {url}")
async def main():
# Limit to 3 concurrent downloads
sem = asyncio.Semaphore(3)
urls = [f"https://example.com/file{i}" for i in range(10)]
await asyncio.gather(*[download(sem, url) for url in urls])
asyncio.run(main())
Event
import asyncio
async def waiter(event, name):
print(f"{name} waiting for event")
await event.wait()
print(f"{name} triggered!")
async def setter(event):
await asyncio.sleep(2)
print("Setting event")
event.set()
async def main():
event = asyncio.Event()
await asyncio.gather(
waiter(event, "Waiter 1"),
waiter(event, "Waiter 2"),
waiter(event, "Waiter 3"),
setter(event)
)
asyncio.run(main())
Condition
import asyncio
async def consumer(condition, items):
async with condition:
await condition.wait()
print(f"Consumed: {items.pop()}")
async def producer(condition, items):
await asyncio.sleep(1)
async with condition:
items.append("item")
print("Produced item")
condition.notify()
async def main():
condition = asyncio.Condition()
items = []
await asyncio.gather(
consumer(condition, items),
producer(condition, items)
)
asyncio.run(main())
Queue
import asyncio
async def producer(queue, n):
for i in range(n):
await asyncio.sleep(0.5)
await queue.put(f"Item {i}")
print(f"Produced: Item {i}")
async def consumer(queue, name):
while True:
item = await queue.get()
print(f"{name} consumed: {item}")
await asyncio.sleep(1)
queue.task_done()
async def main():
queue = asyncio.Queue(maxsize=5)
# Create producers and consumers
producers = [asyncio.create_task(producer(queue, 10))]
consumers = [
asyncio.create_task(consumer(queue, f"Consumer-{i}"))
for i in range(3)
]
# Wait for producers to finish
await asyncio.gather(*producers)
# Wait for queue to be empty
await queue.join()
# Cancel consumers
for c in consumers:
c.cancel()
asyncio.run(main())
Priority Queue
import asyncio
from dataclasses import dataclass, field
from typing import Any
@dataclass(order=True)
class PrioritisedItem:
priority: int
item: Any = field(compare=False)
async def main():
queue = asyncio.PriorityQueue()
# Add items with priorities (lower number = higher priority)
await queue.put(PrioritisedItem(5, "Low priority"))
await queue.put(PrioritisedItem(1, "High priority"))
await queue.put(PrioritisedItem(3, "Medium priority"))
# Retrieve in priority order
while not queue.empty():
item = await queue.get()
print(f"Priority {item.priority}: {item.item}")
asyncio.run(main())
Debugging and Profiling
Debug Mode
import asyncio
import warnings
# Enable debug mode
asyncio.run(main(), debug=True)
# Or set environment variable
# PYTHONASYNCIODEBUG=1 python script.py
# Enable warnings for unawaited coroutines
warnings.simplefilter('always', ResourceWarning)
# In debug mode, asyncio checks:
# - Coroutines not awaited
# - Callbacks taking longer than 100ms
# - Unawaited tasks being destroyed
Logging
import asyncio
import logging
# Enable asyncio logging
logging.basicConfig(level=logging.DEBUG)
# Enable debug mode via asyncio.run(main(), debug=True), or inside a
# coroutine: asyncio.get_running_loop().set_debug(True)
# Custom logger
logger = logging.getLogger('asyncio')
logger.setLevel(logging.DEBUG)
async def main():
logger.info("Starting async operation")
await asyncio.sleep(1)
logger.info("Completed")
asyncio.run(main())
Slow Callback Detection
import asyncio
async def main():
loop = asyncio.get_running_loop()
# Set slow callback threshold (default: 100ms in debug mode)
loop.slow_callback_duration = 0.05 # 50ms
# This callback will be logged as slow
def slow_callback():
import time
time.sleep(0.1) # Blocks for 100ms
loop.call_soon(slow_callback)
await asyncio.sleep(0.2)
asyncio.run(main(), debug=True)
Profiling with cProfile
import asyncio
import cProfile
import pstats
async def cpu_intensive():
total = 0
for i in range(1000000):
total += i
return total
async def main():
results = await asyncio.gather(*[cpu_intensive() for _ in range(10)])
return results
# Profile async code
profiler = cProfile.Profile()
profiler.enable()
asyncio.run(main())
profiler.disable()
# Print stats
stats = pstats.Stats(profiler)
stats.sort_stats('cumulative')
stats.print_stats(10)
Tracing Task Creation
import asyncio
import traceback
async def task_with_error():
await asyncio.sleep(1)
raise ValueError("Something went wrong")
async def main():
loop = asyncio.get_running_loop()
# Custom exception handler
def exception_handler(loop, context):
exception = context.get('exception')
message = context.get('message')
task = context.get('task')
print(f"Exception in task {task}: {message}")
print(f"Exception: {exception}")
# Print task creation traceback. Note: _source_traceback is a
# private, unsupported attribute, only populated in debug mode
if getattr(task, '_source_traceback', None):
print("Task created at:")
traceback.print_stack(task._source_traceback)
loop.set_exception_handler(exception_handler)
task = asyncio.create_task(task_with_error())
await asyncio.sleep(2)
asyncio.run(main(), debug=True)
Performance Monitoring
import asyncio
import time
class TaskMonitor:
def __init__(self):
self.tasks = {}
def track_task(self, task, name):
self.tasks[task] = {
'name': name,
'start': time.time(),
'done': False
}
task.add_done_callback(lambda t: self.on_done(t))
def on_done(self, task):
if task in self.tasks:
info = self.tasks[task]
duration = time.time() - info['start']
print(f"Task '{info['name']}' completed in {duration:.2f}s")
info['done'] = True
async def slow_task(n):
await asyncio.sleep(n)
return f"Completed after {n}s"
async def main():
monitor = TaskMonitor()
tasks = []
for i in range(5):
task = asyncio.create_task(slow_task(i))
monitor.track_task(task, f"Task-{i}")
tasks.append(task)
await asyncio.gather(*tasks)
asyncio.run(main())
Memory Debugging
import asyncio
import gc
import tracemalloc
async def memory_intensive():
data = [list(range(1000)) for _ in range(1000)]
await asyncio.sleep(1)
return len(data)
async def main():
# Start tracing memory allocations
tracemalloc.start()
snapshot1 = tracemalloc.take_snapshot()
await asyncio.gather(*[memory_intensive() for _ in range(10)])
snapshot2 = tracemalloc.take_snapshot()
# Compare snapshots
top_stats = snapshot2.compare_to(snapshot1, 'lineno')
print("Top 10 memory allocation differences:")
for stat in top_stats[:10]:
print(stat)
tracemalloc.stop()
asyncio.run(main())
Quick Reference
Core Functions
| Function | Purpose |
|---|---|
asyncio.run(coro) |
Run a coroutine (main entry point) |
await coro |
Wait for coroutine to complete |
asyncio.create_task(coro) |
Schedule coroutine as task |
asyncio.gather(*coros) |
Run coroutines concurrently, wait for all |
asyncio.wait(tasks) |
Wait for tasks with fine-grained control |
asyncio.wait_for(coro, timeout) |
Wait with timeout |
asyncio.sleep(seconds) |
Async sleep (yields control) |
Task Management
| Method | Purpose |
|---|---|
task.cancel() |
Cancel a task |
task.cancelled() |
Check if task was cancelled |
task.done() |
Check if task completed |
task.result() |
Get task result (or raise exception) |
task.exception() |
Get exception raised by task |
task.add_done_callback(fn) |
Add completion callback |
Synchronisation
| Primitive | Use Case |
|---|---|
asyncio.Lock() |
Mutual exclusion |
asyncio.Semaphore(n) |
Limit concurrent access to n |
asyncio.Event() |
Signal between coroutines |
asyncio.Condition() |
Wait for condition to be true |
asyncio.Queue() |
Producer-consumer pattern |
asyncio.PriorityQueue() |
Priority-based queue |
Network I/O
| Function | Purpose |
|---|---|
asyncio.start_server(callback, host, port) |
Start TCP server |
asyncio.open_connection(host, port) |
Connect TCP client |
loop.create_datagram_endpoint() |
Create UDP endpoint |
loop.run_in_executor(executor, func) |
Run blocking code in executor |
Debugging
| Setting | Purpose |
|---|---|
asyncio.run(main(), debug=True) |
Enable debug mode |
loop.set_debug(True) |
Enable debug on existing loop |
PYTHONASYNCIODEBUG=1 |
Environment variable for debug mode |
loop.slow_callback_duration |
Set threshold for slow callbacks |
loop.set_exception_handler(handler) |
Custom exception handling |
Common Issues and Solutions
Issue: Coroutine Not Awaited
# Problem
async def fetch_data():
return "data"
def main():
result = fetch_data() # ⚠️ RuntimeWarning: coroutine was never awaited
print(result)
# Solution
async def main():
result = await fetch_data() # ✅ Properly awaited
print(result)
asyncio.run(main())
Issue: Blocking I/O in Async Code
import asyncio
import time
# Problem: Blocking the event loop
async def bad_sleep():
time.sleep(5) # ⚠️ Blocks entire event loop
return "done"
# Solution: Use async alternative
async def good_sleep():
await asyncio.sleep(5) # ✅ Yields control to event loop
return "done"
# For unavoidable blocking operations
async def blocking_operation():
loop = asyncio.get_running_loop()
result = await loop.run_in_executor(
None, # Use default ThreadPoolExecutor
time.sleep, 5 # Blocking function
)
return result
Issue: Task Cancellation Not Handled
import asyncio
# Problem: Cancellation not properly handled
async def bad_task():
await asyncio.sleep(10)
return "completed"
# Solution: Handle CancelledError
async def good_task():
try:
await asyncio.sleep(10)
return "completed"
except asyncio.CancelledError:
# Cleanup code
print("Task cancelled, cleaning up...")
raise # Re-raise to mark task as cancelled
async def main():
task = asyncio.create_task(good_task())
await asyncio.sleep(1)
task.cancel()
try:
await task
except asyncio.CancelledError:
print("Task was cancelled")
Issue: TimeoutError Not Caught
import asyncio
async def slow_operation():
await asyncio.sleep(10)
return "done"
# Problem: Timeout not handled
async def bad_timeout():
result = await asyncio.wait_for(slow_operation(), timeout=2.0)
return result # ⚠️ TimeoutError raised
# Solution: Catch TimeoutError
async def good_timeout():
try:
result = await asyncio.wait_for(slow_operation(), timeout=2.0)
return result
except asyncio.TimeoutError:
print("Operation timed out")
return None # Or handle appropriately
asyncio.run(good_timeout())
Issue: Event Loop Already Running
import asyncio
# Problem: Nested asyncio.run() calls
async def inner():
asyncio.run(some_coro()) # ⚠️ RuntimeError: asyncio.run() cannot be called from a running event loop
# Solution 1: Use await instead
async def good_inner():
await some_coro() # ✅ Use await in async context
# Solution 2: Use nest_asyncio for Jupyter/interactive environments
import nest_asyncio
nest_asyncio.apply()
asyncio.run(some_coro()) # Now works in nested contexts
Issue: Memory Leaks from Uncompleted Tasks
import asyncio
# Problem: Tasks created but never awaited
async def bad_pattern():
for i in range(1000):
asyncio.create_task(some_coro()) # ⚠️ Tasks accumulate
# Tasks may be destroyed before completion
# Solution: Track and await tasks
async def good_pattern():
tasks = []
for i in range(1000):
task = asyncio.create_task(some_coro())
tasks.append(task)
# Wait for all tasks
await asyncio.gather(*tasks, return_exceptions=True)
# Solution 2: Use TaskGroup (Python 3.11+)
async def best_pattern():
async with asyncio.TaskGroup() as group:
for i in range(1000):
group.create_task(some_coro())
# All tasks automatically awaited
Issue: Mixing Sync and Async Code
import asyncio
# Problem: Calling async function from sync code
def sync_function():
result = await async_function() # ⚠️ SyntaxError: await outside async function
# Solution 1: Make the function async
async def async_function_wrapper():
result = await async_function() # ✅
return result
# Solution 2: Use asyncio.run() for entry point
def sync_entry():
result = asyncio.run(async_function()) # ✅
return result
# Solution 3: Use run_until_complete (lower level)
def sync_entry_alt():
loop = asyncio.new_event_loop()
try:
result = loop.run_until_complete(async_function())
return result
finally:
loop.close()
Issue: Race Conditions in Concurrent Tasks
import asyncio
# Problem: Race condition
counter = 0
async def bad_increment():
global counter
for _ in range(1000):
temp = counter # ⚠️ Not atomic
await asyncio.sleep(0)
counter = temp + 1
# Solution: Use Lock
async def good_increment(lock):
global counter
for _ in range(1000):
async with lock: # ✅ Atomic operation
counter += 1
async def main():
lock = asyncio.Lock()
await asyncio.gather(*[good_increment(lock) for _ in range(10)])
print(f"Counter: {counter}") # Correct value
asyncio.run(main())
Issue: Exception Handling in gather()
import asyncio
async def failing_task():
await asyncio.sleep(1)
raise ValueError("Task failed")
async def working_task():
await asyncio.sleep(2)
return "success"
# Problem: Exception stops all tasks
async def bad_gather():
results = await asyncio.gather(
failing_task(),
working_task()
) # ⚠️ Raises ValueError, working_task may not complete
# Solution: Use return_exceptions=True
async def good_gather():
results = await asyncio.gather(
failing_task(),
working_task(),
return_exceptions=True # ✅ Returns exceptions as values
)
for i, result in enumerate(results):
if isinstance(result, Exception):
print(f"Task {i} failed: {result}")
else:
print(f"Task {i} succeeded: {result}")
asyncio.run(good_gather())
Issue: Deadlock with Synchronisation Primitives
import asyncio
# Problem: Deadlock scenario
async def bad_pattern(lock1, lock2):
async with lock1:
await asyncio.sleep(0.1)
async with lock2: # ⚠️ Potential deadlock
pass
# Solution: Always acquire locks in same order
async def good_pattern(lock1, lock2):
locks = sorted([lock1, lock2], key=id) # ✅ Consistent ordering
async with locks[0]:
async with locks[1]:
pass
# Alternative: Use timeout
async def pattern_with_timeout(lock1, lock2):
async with lock1:
try:
async with asyncio.timeout(5.0): # Python 3.11+
async with lock2:
pass
except asyncio.TimeoutError:
print("Lock acquisition timed out")
Issue: Queue Deadlock
import asyncio
# Problem: Queue deadlock
async def bad_queue_pattern():
queue = asyncio.Queue(maxsize=1)
# Producer
await queue.put("item1")
await queue.put("item2") # ⚠️ Blocks forever if no consumer
# Solution: Use proper producer-consumer pattern
async def producer(queue):
for i in range(10):
await queue.put(f"item{i}")
await asyncio.sleep(0.1)
async def consumer(queue):
while True:
item = await queue.get()
print(f"Consumed: {item}")
queue.task_done()
await asyncio.sleep(0.5)
async def good_queue_pattern():
queue = asyncio.Queue()
producer_task = asyncio.create_task(producer(queue))
consumer_task = asyncio.create_task(consumer(queue))
await producer_task
await queue.join() # Wait for all items to be processed
consumer_task.cancel()
asyncio.run(good_queue_pattern())