Why must you await asyncio's StreamWriter.drain() after calling write()?
answer
- write() returns before anything is sent
- The buffer it fills lives in your heap
- Two thresholds on that buffer
- TCP pushes back only on the kernel
- It also raises when the peer is gone
basics
~20 sStreamWriter.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.
solid answer
~50 s`writer.write(data)` never blocks: it appends to the transport's in-process write buffer and the event loop flushes it when the socket is writable. If the peer reads slower than you write, nothing pushes back, and the buffer — plain heap memory — grows until the process dies. `await writer.drain()` is asyncio's backpressure hook: it returns immediately while the buffered size is under the transport's **low-water mark**, and suspends the calling coroutine once the buffer has passed the **high-water mark** (64 KiB by default, with the low mark at a quarter of it), resuming only when the transport drains back down. It also re-raises a connection error such as a reset peer, which a bare `write()` would silently swallow. Note what it is *not*: it does not flush to the OS and it is no delivery acknowledgement.
code
python · 18 linesimport asyncio
async def stalled_consumer(reader, writer):
await asyncio.sleep(0.5)
writer.close()
async def main():
server = await asyncio.start_server(stalled_consumer, "127.0.0.1", 0)
port = server.sockets[0].getsockname()[1]
async with server:
_, writer = await asyncio.open_connection("127.0.0.1", port)
writer.transport.set_write_buffer_limits(high=16 * 1024)
for _ in range(200):
writer.write(b"x" * 8192)
print("bytes queued in memory:", writer.transport.get_write_buffer_size())
writer.close()
asyncio.run(main())go deeper
Remember the idiom: write() then await drain(). Know that write() is not awaited and that skipping drain() in a loop is the mistake, even if you cannot yet explain the water marks.
Explain the mechanism: the transport's write buffer, the high- and low-water marks, and that drain() suspends only once the high mark is crossed. Say plainly that drain() is not a flush.
Diagnose it in production: rising RSS on a service with a slow downstream, writes that never raise into a dead peer, and the choice between tuning set_write_buffer_limits and putting a bounded asyncio.Queue in front with an explicit shedding policy.
Own the end-to-end backpressure story: where the system is allowed to buffer, how much memory that costs per connection at fleet scale, and whether the right answer is to slow the producer, shed load, or move the buffer to durable storage.
## The problem drain() exists to solve Consider a metrics scraper that pulls samples from a fleet of agents and streams them onward to a collector over a single TCP connection. Producing a sample is cheap; the collector is slow, and after a restart it takes a 45-second cold start before it reads at full rate. The scraper's forwarding coroutine does `writer.write(sample)` in a loop. `asyncio.StreamWriter.write()` is a plain synchronous method. It calls the transport's `write`, which tries a non-blocking send on the socket; whatever the kernel will not take right now is appended to a `bytearray` the transport owns in your process's heap, and the loop registers interest in socket-writability to flush the rest later. `write()` then returns. It cannot signal "slow down" — it has no return value you check and it raises nothing. So during those 45 seconds nothing stops the loop. The socket send buffer fills, the peer's receive window closes, TCP's own backpressure reaches the kernel — and stops there. Above it, Python keeps appending. RSS climbs until the container is OOM-killed. TCP flow control protected the network and did nothing at all for your process. ## What drain() actually does Every asyncio transport has a **write buffer** with two thresholds, the high-water and low-water marks, settable with `set_write_buffer_limits`. By default the high mark is 64 KiB and the low mark is a quarter of it, 16 KiB. When the buffered size crosses the high mark the transport calls its protocol's `pause_writing()`; when it falls back below the low mark it calls `resume_writing()`. The streams layer turns those two callbacks into a future. `await writer.drain()` awaits that future. Concretely: * buffer below the low-water mark → `drain()` returns essentially immediately (it still yields to the loop in some paths, which is itself useful — it gives other tasks a turn); * buffer above the high-water mark → the coroutine suspends until the transport has flushed enough bytes to the socket to get back under the low mark. The result is exactly the coupling you want: the producer's speed becomes a function of the consumer's speed, and memory use is bounded by the high-water mark rather than by how fast you can generate data. The correct idiom is `write(...)` then `await drain()` — write as many chunks as you like, then drain once per batch. Draining after every tiny write is not wrong, only slower, since it costs a loop turn. ## The second job: surfacing errors `drain()` also checks the connection's state. If the peer has reset the connection or the transport is closing, `drain()` raises — typically `ConnectionResetError` or `BrokenPipeError`. A `write()` into a dead connection raises nothing, so a loop that never drains can write happily into the void long after the peer has gone. In a scraper this is the off-by-one boundary case that bites: the last batch before a collector restart appears to have been sent, no exception is raised anywhere, and the samples are silently lost. Draining is where you find out. ## What drain() does *not* mean Three misconceptions, all common: 1. **It is not a flush.** `drain()` returning does not mean the bytes reached the socket, let alone the peer. It means the in-process buffer is small again. 2. **It is not an acknowledgement.** No application-level delivery guarantee comes out of it. If you need one, the peer must send something back and you must read it. 3. **It is not per-write mandatory.** Nothing enforces draining. Code without it works fine in tests, where the consumer is fast, and fails in production, where it is not. ## Tuning and the alternatives `writer.transport.set_write_buffer_limits(high=..., low=...)` moves the marks. Raising the high mark buys tolerance for bursts at the cost of memory; lowering it makes the producer feel backpressure sooner and keeps latency predictable. If you cannot afford to suspend the producer — say the samples are being generated by something you do not control — the shape to reach for is a bounded `asyncio.Queue` in front of the writer, with an explicit policy for what happens when it is full: drop the oldest sample, drop the newest, or block the upstream. That is a decision to make deliberately, and the queue's `maxsize` makes it visible in the code. Silently unbounded buffering is the one option that is never right. ## The one-line version `write()` is fire-and-forget into unbounded process memory; `drain()` is where TCP's backpressure is finally allowed to reach your coroutine, and where connection errors surface.
- Does awaiting drain() guarantee the peer received the bytes?No. It guarantees only that the transport's in-process write buffer has fallen back below its low-water mark. The bytes may still be in the kernel send buffer or in flight, and the peer may never read them. Delivery is an application-level fact: if you need it, the peer has to acknowledge and you have to read that acknowledgement.
- You cannot slow the producer down — it is generating samples on a fixed schedule. What do you do instead of drain()?Put a bounded `asyncio.Queue` between the producer and a single writer coroutine that does write-then-drain. The queue's `maxsize` makes the buffer explicit, and you choose the overflow policy out loud: shed the oldest samples, shed the newest, or let `put()` block and push back on the producer. The failure mode you are removing is the unbounded transport buffer, not the buffering itself.
- How do you change how much a connection buffers before drain() suspends the writer?`writer.transport.set_write_buffer_limits(high=..., low=...)`. The default high-water mark is 64 KiB with the low mark at a quarter of it. A higher mark absorbs bursts at the cost of resident memory per connection — which multiplies by connection count — while a lower mark applies backpressure sooner and keeps latency and memory tighter.
saying these in an interview costs you the question
- Says drain() flushes the data to the peer
- Believes write() blocks when the socket is full
- Thinks TCP flow control already bounds the process's memory
- Claims drain() confirms the peer received the bytes
- Assumes the transport raises when the write buffer overflows