asyncio Streams, Timeouts, Semaphores & Backpressure
Building enterprise-grade asynchronous networking systems requires managing network streams (StreamReader / StreamWriter), enforcing strict operation deadlines (asyncio.timeout), rate-limiting concurrent tasks (asyncio.Semaphore), and managing Backpressure using bounded queues (asyncio.Queue).
This chapter details High-Level Async Streams, asyncio.timeout (Python 3.11+), async concurrency control with Semaphore, and memory-safe Backpressure queues.
1. High-Level Async Network Streams (StreamReader / StreamWriter)
asyncio provides high-level socket stream abstractions over low-level socket file descriptors:
import asyncio
async def fetch_http_header(host: str) -> str:
# Open non-blocking TCP socket stream
reader, writer = await asyncio.open_connection(host, 80)
# Write HTTP GET request
request = f"GET / HTTP/1.1\r\nHost: {host}\r\nConnection: close\r\n\r\n"
writer.write(request.encode("utf-8"))
await writer.drain() # Flushes buffer and yields control until write buffer drains!
# Read line asynchronously
line = await reader.readline()
writer.close()
await writer.wait_closed()
return line.decode("utf-8").strip()The writer.drain() Invariant:
Calling writer.write(data) buffers data in CPython memory. Calling await writer.drain() pauses the coroutine if the OS write buffer is full, resuming execution only when the underlying socket buffer has drained.
2. Asynchronous Timeouts (asyncio.timeout in Python 3.11+)
Prior to Python 3.11, timeouts were applied using asyncio.wait_for(). Python 3.11 introduced the cleaner context-manager-based asyncio.timeout():
import asyncio
async def fetch_with_deadline(url: str):
# Enforces a strict 2.5-second deadline for the context block!
try:
async with asyncio.timeout(2.5):
response = await fetch_remote_api(url)
return response
except TimeoutError:
print(f"Request to {url} timed out after 2.5 seconds!")If the code inside the async with asyncio.timeout() block exceeds the specified deadline, asyncio automatically cancels the underlying task and raises a TimeoutError.
3. Rate-Limiting Concurrency with asyncio.Semaphore
Unrestricted asyncio.gather(*[fetch(id) for id in range(10000)]) spawns 10,000 concurrent HTTP requests simultaneously, overwhelming target servers and triggering socket connection refusals (Too Many Open Files).
Use asyncio.Semaphore to bound concurrent task execution:
import asyncio
# Allow at most 10 concurrent requests simultaneously!
sem = asyncio.Semaphore(10)
async def safe_fetch(user_id: int):
async with sem: # Acquires semaphore token; waits if 10 tasks are active
return await fetch_user_api(user_id)
async def main():
tasks = [safe_fetch(i) for i in range(1000)]
results = await asyncio.gather(*tasks)4. Async Backpressure Queues (asyncio.Queue)
In asynchronous producer-consumer pipelines, if producers generate events faster than consumers can process them, un-bounded queues grow infinitely, causing OOM memory crashes.
Backpressure is the mechanism where consumers signal producers to slow down:
# Create bounded queue with maxsize=100
queue = asyncio.Queue(maxsize=100)
async def producer():
for i in range(1000):
# Await put(): Blocks producer when queue size hits 100!
await queue.put(f"item_{i}")
async def consumer():
while True:
item = await queue.get()
await process(item)
queue.task_done()