skip to content

Streams and Protocols

Talking to sockets without threads: open_connection and start_server hand you a StreamReader and StreamWriter, with transports and protocols underneath. Interviewers ask where the backpressure lives.

part ofPythonoverview, primer and where to startread it →
on this pageshow

questions

4

What does asyncio.open_connection return, and what do you do with each object?

level: juniorimportance: must knowfreq 55%

answer

  1. Awaiting it hands back two objects
  2. One pulls bytes, one pushes them
  3. Only one side is awaited
  4. Shutdown takes two calls, not one
  5. start_server is the mirror image

basics

~10 s

Awaiting asyncio.open_connection returns a pair: a StreamReader you await bytes from, and a StreamWriter you push bytes into. Reads are awaited, writes are not, and you finish with close() followed by an awaited wait_closed().

solid answer

~40 s

`await asyncio.open_connection(host, port)` returns a two-tuple `(reader, writer)`. The `StreamReader` is the pull side: `await reader.readline()`, `await reader.readexactly(n)` or `await reader.read(n)`. The `StreamWriter` is the push side, and `writer.write(data)` is a **plain method, not a coroutine** — it hands bytes to the underlying transport and returns immediately, so you follow a burst of writes with `await writer.drain()` to respect backpressure. Shutdown is two steps: `writer.close()` starts the close, and `await writer.wait_closed()` waits for it to finish. The server mirror is `asyncio.start_server(handler, host, port)`, which returns an `asyncio.Server` and calls your handler coroutine with a fresh `(reader, writer)` pair per accepted connection; `async with server:` plus `await server.serve_forever()` is the usual run shape.

code

python · 21 lines
python
import asyncio

async def handle(reader, writer):
    line = await reader.readline()
    writer.write(line.upper())
    await writer.drain()
    writer.close()
    await writer.wait_closed()

async def main():
    server = await asyncio.start_server(handle, "127.0.0.1", 0)
    port = server.sockets[0].getsockname()[1]
    async with server:
        reader, writer = await asyncio.open_connection("127.0.0.1", port)
        writer.write(b"cpu_seconds 12\n")
        await writer.drain()
        print(await reader.readline())
        writer.close()
        await writer.wait_closed()

asyncio.run(main())

go deeper

for a junior

Recall the shape: awaiting open_connection gives a reader and a writer, reads are awaited, writes are not, and closing is close() then await wait_closed(). Be able to write a ten-line client from memory.

for a middle

Explain why the two sides are asymmetric, what the server callback gets per connection, and how you detect EOF. Show a handler that closes its writer in a finally block.

for a senior

Demonstrate that you have operated this: timeouts around reads, framing instead of raw read(n), guaranteed close paths, and what a skipped wait_closed() looks like in a process that runs for weeks.

for a principal

Own the call of whether a service should speak a raw stream protocol at all versus adopting an existing framed protocol library, and what that choice costs in interoperability, observability and on-call load.

## The two objects `asyncio.open_connection(host, port)` is a coroutine function. Awaiting it opens a TCP connection and gives you back exactly two objects: * an `asyncio.StreamReader` — the **read** side. Everything on it is a coroutine, because bytes may not have arrived yet: `await reader.read(n)`, `await reader.readline()`, `await reader.readexactly(n)`, `await reader.readuntil(sep)`. * an `asyncio.StreamWriter` — the **write** side. `writer.write(data)` and `writer.writelines(chunks)` are ordinary synchronous methods; only `await writer.drain()` and `await writer.wait_closed()` are coroutines. That asymmetry is the single most useful thing to remember. Reading genuinely has to wait for the peer, so it suspends. Writing never waits: the bytes are appended to the transport's write buffer and the event loop flushes them to the socket when the socket becomes writable. `write()` therefore cannot fail with "would block" and cannot apply backpressure on its own — that is what `drain()` is for. ## The server side `asyncio.start_server(client_connected_cb, host, port)` is the mirror image. It returns an `asyncio.Server`. Every time a client connects, asyncio creates a fresh `StreamReader`/ `StreamWriter` pair and schedules your callback with them as its two arguments. The callback is normally an `async def` coroutine function; asyncio wraps each invocation in its own task, so every connection is handled concurrently without you creating threads. A common misreading is that `start_server` blocks. It does not — awaiting it only binds the listening socket and starts accepting. To keep the process alive you either `await server.serve_forever()` or keep the enclosing coroutine alive some other way. `async with server:` closes the listening socket on exit. ## Closing correctly `writer.close()` is synchronous and only *requests* the close: it tells the transport to flush what it can and then shut the socket down. The connection is not gone when it returns. `await writer.wait_closed()` is what actually waits for the underlying transport to finish closing and re-raises any error that happened on the way out. Skipping it in a short-lived script usually "works", which is exactly why it becomes a habit — and then in a long-running process you get sockets in a half-closed state, `ResourceWarning`s at interpreter shutdown, and data written just before the close that never left the buffer. On the server object the same two-step applies: `server.close()` stops accepting new connections, and `await server.wait_closed()` waits. Python 3.13 added `asyncio.Server.close_clients()` and `asyncio.Server.abort_clients()` for the case where you also want to tear down the connections that are already open rather than waiting for handlers to finish on their own. ## What is underneath Streams are a façade. Underneath sits asyncio's transport/protocol layer: a transport owns the socket and the write buffer, and a protocol receives callbacks. `open_connection` builds a protocol that pushes incoming bytes into the `StreamReader`'s internal buffer and wakes whichever coroutine is awaiting a read. You rarely need to know that to use streams, but it explains why a `StreamReader` has no `write` and a `StreamWriter` has no `read`: they are two views on one connection, not two connections. ## Errors and timeouts belong to the caller Neither `open_connection` nor the reader and writer carry a timeout of their own. A read against a peer that connected and then went silent waits forever, which in a long-running process is indistinguishable from a hang. Wrapping the read in `async with asyncio.timeout(...)` is the caller's job, as is deciding what a timeout means for that connection — usually closing it, because a half-consumed stream cannot be reused. Connection failures surface where you would expect: `open_connection` raises `OSError` and its subclasses such as `ConnectionRefusedError`, and a mid-stream failure raises `ConnectionResetError` from a read or from `await writer.drain()`. ## Typical shape A client is usually: connect, write a request, `drain`, read a framed response, close, wait closed — wrapped in `try/finally` so the close happens even when the read raises. A server handler is usually a loop reading one framed request at a time until the reader hits EOF (a read returning `b""`), writing a reply per request, and closing in a `finally`. ## What interviewers listen for Three things, in order: that you know `write()` is not awaited; that you know reads must be *framed* (a `read(n)` can return fewer than `n` bytes, so you use `readline`/`readexactly`/ `readuntil` when you need a message boundary); and that close is `close()` **then** `await wait_closed()`. A candidate who awaits `write()` or who thinks `close()` is sufficient has not run this code against a real peer.

  • Why is StreamWriter.write() a plain method while StreamReader.read() is a coroutine?
    Reading may have to wait for bytes that have not arrived, so it must be able to suspend. Writing never waits: `write()` appends to the transport's write buffer and returns, and the event loop flushes it when the socket is writable. The waiting was factored out into `await writer.drain()` so that the common case — write and move on — costs nothing.
  • What does the callback passed to asyncio.start_server receive, and what happens if it raises?
    It is called with a fresh `(reader, writer)` pair for each accepted connection, normally as an `async def` coroutine function that asyncio wraps in its own task. If it raises, the exception is reported through the loop's exception handler and that connection's handling stops; the server keeps accepting. So a handler should own its own `try/finally` for closing the writer rather than relying on the server to clean up.
  • How do you know a peer has stopped sending on an asyncio StreamReader?
    A read returns an empty bytes object: `await reader.read(n)` giving `b""`, or `await reader.readline()` giving `b""` with nothing buffered, means EOF. `reader.at_eof()` reports the same condition. Never treat `b""` as "nothing right now" — a stream read that has nothing yet suspends instead of returning empty.

saying these in an interview costs you the question

  • Says StreamWriter.write() must be awaited
  • Treats writer.close() as enough, never awaits wait_closed()
  • Thinks open_connection returns a raw socket object
  • Assumes read(n) always returns exactly n bytes
  • Believes awaiting start_server blocks until the server stops

context

open as a page

Why must you await asyncio's StreamWriter.drain() after calling write()?

level: middleimportance: must knowfreq 60%

basics

~20 s

StreamWriter.write() only queues bytes in the transport's write buffer and returns at once, so a fast producer against a slow peer grows that buffer without bound. Awaiting drain() suspends the producing coroutine until the buffer drops back below its low-water mark.

open as a page

How do asyncio's StreamReader.read, readline and readexactly differ at end of stream?

level: middleimportance: should knowfreq 42%

basics

~20 s

read returns whatever is available, and empty bytes at EOF. readline returns a partial line without its newline at EOF, never raising. readexactly and readuntil raise asyncio.IncompleteReadError when the stream ends early, carrying the bytes already read.

open as a page

What is asyncio's transport and protocol layer, and when would you use it instead of streams?

level: seniorimportance: nice to knowfreq 24%

basics

~20 s

Transports own the socket and its write buffer; protocols are your callback object, receiving connection_made, data_received, eof_received and connection_lost. Streams are a coroutine-friendly façade built on that pair. Drop to protocols for lower per-connection overhead and zero-copy reads.

open as a page