Reducing Copies with Buffered Protocols¶
asyncio offers three levels for receiving data: streams (StreamReader), Protocol with data_received(bytes), and BufferedProtocol, where the event loop reads directly into a buffer you provide. Each level down removes copies and Python-level work, and adds code you must get right. Measured on Python 3.14 receiving 2 GiB over loopback: streams read at 2,248 MiB/s, a plain Protocol at 7,812 MiB/s, and BufferedProtocol at 6,912–8,140 MiB/s — for raw bytes, leaving streams is the big step and the buffered protocol adds little. Parsing small messages is different. With 2 million 200-byte length-prefixed frames, StreamReader.readexactly handled 1.26 million frames per second, a Protocol that accumulated into a bytearray and sliced out each frame 3.29 million, and a BufferedProtocol that parsed frames as zero-copy memoryviews in its own buffer 6.02 million. This guide shows when the buffered protocol is worth it and how to write one correctly.
Prerequisites¶
- Python 3.11+, stdlib only.
- Transports and protocols, from Streams, Transports & Protocols.
- A profile showing receive-side parsing as a significant cost; otherwise stay with streams.
1. Know where the copies are¶
Each level hands data to your code differently:
# Streams: the transport calls data_received on an internal protocol, which appends to
# StreamReader's bytearray; readexactly() then copies the bytes out again.
frame = await reader.readexactly(n)
# Protocol: you receive a new bytes object per socket read and manage your own buffer.
def data_received(self, data: bytes) -> None:
self.buf += data # copy 1: into your accumulator
frame = bytes(self.buf[4:4 + n]) # copy 2: out of it
# BufferedProtocol: the loop calls recv_into() on a buffer you own; frames can be
# memoryview slices of that buffer with no copy at all.
def get_buffer(self, sizehint: int) -> memoryview: ...
def buffer_updated(self, nbytes: int) -> None: ...
For bulk transfer the copies are cheap relative to the system calls, which is why raw throughput was similar for Protocol and BufferedProtocol. For many small messages, the per-message Python work dominates: an await and a buffer copy per readexactly call for streams, a slice and a bytes object per frame for Protocol. Removing those is where the measured 4.8× over streams comes from.
Verify: profile the receive path; if readexactly, bytearray operations or bytes construction are near the top, a lower level will help.
2. Write a BufferedProtocol with a reusable buffer¶
The loop asks for a buffer with get_buffer, reads into it, then reports how many bytes it wrote with buffer_updated. Keep one large buffer and track the unparsed region:
import asyncio
import struct
class FrameProtocol(asyncio.BufferedProtocol):
def __init__(self, on_frame, size: int = 1 << 20) -> None:
self.on_frame = on_frame
self.buf = bytearray(size)
self.view = memoryview(self.buf)
self.start = 0 # first unparsed byte
self.end = 0 # one past the last received byte
def get_buffer(self, sizehint: int) -> memoryview:
if self.start: # move a partial frame to the front
remaining = self.end - self.start
self.view[:remaining] = self.view[self.start:self.end]
self.start, self.end = 0, remaining
if self.end == len(self.buf):
raise BufferError("frame larger than buffer")
return self.view[self.end:]
def buffer_updated(self, nbytes: int) -> None:
self.end += nbytes
pos, end, view = self.start, self.end, self.view
while end - pos >= 4:
(length,) = struct.unpack_from("!I", view, pos)
if end - pos - 4 < length:
break # incomplete frame: wait for more
self.on_frame(view[pos + 4:pos + 4 + length]) # zero-copy view
pos += 4 + length
self.start = pos
Compaction happens once per read, moving only the trailing partial frame, not once per frame. The buffer size is also the maximum frame size; check declared lengths against it and close the connection on oversize frames, rather than letting get_buffer fail. Protocol callbacks are plain functions, not coroutines — they must not block, and they cannot await.
Verify: feed the protocol frames split at every possible byte boundary in a unit test (call get_buffer and buffer_updated directly); every frame is delivered exactly once.
3. Handle the memoryview lifetime¶
The zero-copy frames are views into a buffer that will be overwritten by the next read. Use them immediately, or copy:
def on_frame(frame: memoryview) -> None:
kind = frame[0] # fine: read during the callback
if kind == HEARTBEAT:
return
message = parse(frame) # fine: parse into new objects now
queue.put_nowait(message)
def on_frame_wrong(frame: memoryview) -> None:
queue.put_nowait(frame) # BUG: contents change after the next read
This is the main correctness risk of the buffered approach: a view kept past the callback silently changes under you when the buffer is compacted or refilled. Parsing into Python objects inside the callback — or copying with bytes(frame) when the raw bytes must outlive it — keeps it safe. The parsing work itself then dominates, which is why real-world gains are usually smaller than the 6.02 million frames per second measured for parsing alone.
Verify: a test that holds frames past the callback detects corruption — then remove the bug it was written to catch.
4. Apply backpressure from a protocol¶
Without streams, there is no await reader.read() to slow you down; data arrives as fast as the peer sends. If frames are handed to slower async consumers, pause the transport when your queue is full:
class BackpressuredFrames(FrameProtocol):
def connection_made(self, transport: asyncio.Transport) -> None:
self.transport = transport
def deliver(self, message) -> None:
self.queue.put_nowait(message)
if self.queue.qsize() >= HIGH_WATER and not self.paused:
self.transport.pause_reading() # stop calling get_buffer / recv
self.paused = True
def consumer_took_one(self) -> None: # called by the consumer task
if self.paused and self.queue.qsize() <= LOW_WATER:
self.transport.resume_reading()
self.paused = False
pause_reading() stops the loop from reading the socket; the kernel buffer fills and TCP flow control slows the sender, exactly like a stream reader that is not being read. Use high and low water marks so the transport does not flip on every message. On the sending side, protocols get pause_writing() and resume_writing() callbacks from the transport when its write buffer crosses its own water marks.
Verify: with a deliberately slow consumer, process memory stays bounded and the sender is slowed.
5. Decide whether it is worth it¶
Most services should stay with streams:
# Rough decision based on measured receive rates on one core:
# streams: ~2.2 GiB/s raw, ~1.3 M small frames/s
# Protocol: ~7.8 GiB/s raw, ~3.3 M small frames/s
# BufferedProtocol: ~7-8 GiB/s raw, ~6.0 M small frames/s
# If your service handles 50,000 messages/s, receive parsing is a few % of one core either way.
Streams give you await, timeouts with asyncio.timeout, readable sequential code and backpressure for free. A buffered protocol makes sense for a component whose whole job is moving many small messages — a market-data feed handler, a metrics ingestion endpoint, a proxy, a message broker client — and only after a profile shows receive-side parsing as the bottleneck. Uvicorn's httptools protocol, asyncpg and aiohttp all use protocols internally for this reason, so you get the benefit without writing one whenever you use them.
Verify: the decision is backed by a profile of the real workload, and the protocol has tests at every split boundary.
Verification¶
A buffered protocol is justified and correct when:
- A profile showed receive parsing as the bottleneck before it was written.
- Frames are parsed in place and views are used or copied inside the callback.
- Oversize frames are rejected, and partial frames survive compaction.
- Backpressure pauses reading when consumers fall behind.
Diagnostic Hook: measure frames per second and CPU per frame in the receive path, and count compactions and their byte sizes. Compactions moving large amounts suggest the buffer is too small for the typical burst; CPU per frame that does not fall after moving to a buffered protocol means the cost is in processing, not receiving.
Pitfalls & edge cases¶
- Keeping memoryviews after the callback. Their contents change on the next read.
- Frames larger than the buffer. Check declared lengths before trusting them.
- No backpressure. Protocols receive as fast as the peer sends.
- Optimizing the wrong layer. Raw throughput was similar for Protocol and BufferedProtocol.
Frequently Asked Questions¶
What is asyncio BufferedProtocol?
A protocol where the event loop reads socket data directly into a buffer you provide through get_buffer, then calls buffer_updated with the number of bytes written, avoiding a new bytes object per read.
Is BufferedProtocol faster than Protocol?
For raw bytes, barely: both received about 7-8 GiB/s on loopback in testing. For many small frames parsed in place it was much faster: 6.02 million frames per second against 3.29 million.
Is StreamReader slow for high message rates?
Relative to protocols, yes: it parsed 1.26 million 200-byte frames per second in testing, against 6.02 million with a BufferedProtocol. For most services that rate is far beyond what they need.
How do I apply backpressure in an asyncio Protocol?
Call transport.pause_reading() when your downstream queue reaches a high-water mark and transport.resume_reading() when it drains below a low-water mark.
Related¶
- Streams, Transports & Protocols — up to the topic overview.
- Building a TCP proxy with asyncio streams — where streams are fast enough and much simpler.
- Network I/O & Protocol Handling — the section overview.