Reading stdin Asynchronously in asyncio Scripts¶
Pipeline scripts read records from stdin — cat ids.txt | python enrich.py, kubectl get pods -o name | python check.py — and process them concurrently. Reading stdin with a blocking sys.stdin.readline() freezes the event loop while it waits, so async scripts need another way in, and the obvious options differ by orders of magnitude. Measured on Python 3.14 with a 1,000,000-line, 25.8 MB file: await asyncio.to_thread(sys.stdin.buffer.readline) per line read 38,000 lines/s — 26 seconds for the file; a StreamReader attached with loop.connect_read_pipe read 1.69 million lines/s from a pipe but raised ValueError: Pipe transport is for pipes/sockets only. when stdin was redirected from a file; reading 64 KiB chunks in a thread and splitting them read the whole file in 0.03–0.04 s; and a reader thread feeding batches of 1,000 lines into an asyncio.Queue took 0.07 s. There is a side effect too: on a terminal, connect_read_pipe switched stdin to non-blocking mode and left it that way after the transport closed. This guide reads stdin quickly from any source and leaves the terminal as it found it.
Prerequisites¶
- Python 3.11+; standard library only.
- Thread offloading, from running blocking SDK calls with asyncio.to_thread.
- The topic overview, Async Scripts & CLIs.
1. Do not hop to a thread per line¶
The most common pattern looks harmless:
import asyncio
import sys
async def read_lines_slow():
while line := await asyncio.to_thread(sys.stdin.buffer.readline):
yield line
Measured: 26.4 s for one million lines from a pipe, 25.9 s from a file — 38,000 lines per second. Each line costs a round trip to the default thread pool: submit a job, wake a worker, run readline (which returns immediately from the buffer for all but one call in a few hundred), and schedule the result back onto the loop. The I/O was never the cost; the thread hop was. For an interactive prompt that reads one line at a time this is fine. For a data pipeline it makes stdin the bottleneck of the whole script.
Verify: time your reader alone on a file of realistic size, with the processing removed; it should be far faster than the processing.
2. Read chunks in a thread and split on the loop¶
Amortize the thread hop over many lines by reading a block at a time:
async def read_lines(chunk_size: int = 1 << 16):
rest = b""
while chunk := await asyncio.to_thread(sys.stdin.buffer.read1, chunk_size):
lines = (rest + chunk).split(b"\n")
rest = lines.pop() # the last piece may be an incomplete line
for line in lines:
yield line
if rest:
yield rest # final line without a trailing newline
Measured: one million lines counted in 0.03–0.04 s from either a pipe or a redirected file. read1 returns whatever is available up to the limit, so an interactive or slow producer still gets its lines through promptly instead of waiting for 64 KiB to accumulate. The rest handling matters: a chunk boundary falls in the middle of a line almost every time, and forgetting it silently corrupts records at the boundaries. This approach works for every kind of stdin — pipe, file, terminal, socket — because the blocking read happens in a thread rather than through the event loop's selector.
Verify: feed a file whose last line has no trailing newline and whose lines are longer than the chunk size; every line comes through intact.
3. Use connect_read_pipe only when stdin is a pipe¶
loop.connect_read_pipe attaches stdin to the event loop's selector and gives you a real StreamReader, with readline(), readuntil() and async for:
async def stdin_reader() -> asyncio.StreamReader:
loop = asyncio.get_running_loop()
reader = asyncio.StreamReader(limit=2 ** 16)
await loop.connect_read_pipe(lambda: asyncio.StreamReaderProtocol(reader), sys.stdin)
return reader
async def main() -> None:
reader = await stdin_reader()
async for line in reader:
handle(line)
Measured: 1.69 million lines per second from cat file | python script.py. With python script.py < file it raised ValueError: Pipe transport is for pipes/sockets only., because a regular file cannot be registered with epoll. Scripts are run both ways, so either check the type of stdin first — stat.S_ISFIFO(os.fstat(0).st_mode) — or use the chunked reader, which handles both. The limit argument bounds a single line's length; a longer line raises ValueError from readline, which is a feature when input is untrusted.
Verify: run the script with both cat file | and < file; both work, or the script chooses its reader by stdin type.
4. Restore blocking mode on the terminal¶
When stdin is a terminal, the file description is shared with the shell that launched the script. Measured under a pseudo-terminal:
print(os.get_blocking(0)) # True before connecting
tr, _ = await loop.connect_read_pipe(lambda: asyncio.StreamReaderProtocol(reader), sys.stdin)
print(os.get_blocking(0)) # False while connected
tr.close()
print(os.get_blocking(0)) # False after close, and after the loop exited
asyncio switches the descriptor to non-blocking mode and does not switch it back. The next program to read the same terminal may get BlockingIOError (EAGAIN) instead of waiting for input — a confusing failure in a different process. Restore the flag yourself:
was_blocking = os.get_blocking(0)
try:
await run_with_stdin_reader()
finally:
os.set_blocking(0, was_blocking)
This matters most for interactive tools and for scripts composed in shell pipelines with other programs reading the same terminal. The thread-based readers never touch the flag.
Verify: os.get_blocking(0) is True after the script's main coroutine returns, when run from a terminal.
5. Process lines concurrently with a bound¶
Reading fast only helps if processing keeps up without unbounded memory. Combine the chunked reader with a fixed number of workers and a bounded queue, so a fast producer waits when workers fall behind:
async def main(concurrency: int = 32) -> int:
queue: asyncio.Queue[bytes | None] = asyncio.Queue(maxsize=concurrency * 4)
failures = 0
async def worker() -> None:
nonlocal failures
while (line := await queue.get()) is not None:
try:
await enrich(line.decode())
except Exception:
failures += 1
async with asyncio.TaskGroup() as tg:
for _ in range(concurrency):
tg.create_task(worker())
async for line in read_lines():
await queue.put(line) # waits when the queue is full
for _ in range(concurrency):
await queue.put(None) # one stop signal per worker
return 1 if failures else 0
The bounded queue is the backpressure: when enrich is slow, queue.put suspends the reader, which stops calling read1, which lets the upstream program's pipe buffer fill and block it in turn. Measured with yes as an endless producer, 32 workers and 1 ms of simulated work per line: the queue sat at its limit of 128 items, about 27,500 lines were processed per second, and peak RSS stayed at 20.2 MiB for every one-second sample. Memory stays proportional to the queue size, not to the input size — the same discipline as bounded queues with backpressure. Returning a status from main lets the script report partial failure through its exit code.
Verify: piping an endless producer (yes | python script.py) into a script with slow processing keeps memory flat.
Verification¶
Async stdin handling is correct when:
- No code hops to a thread per line in a data path.
- The reader works with pipes, files and terminals, or chooses by stdin type.
- Blocking mode is restored whenever
connect_read_pipewas used on a terminal. - Processing is bounded by a queue size and a worker count.
Diagnostic Hook: log lines read per second and queue depth every few seconds during a run. A read rate near 40,000 lines per second with an empty queue means the reader is the bottleneck — almost always a thread hop per line; a full queue with a falling read rate means the workers are, which is the healthy case.
Pitfalls & edge cases¶
to_thread(readline)in a loop. Measured: 38,000 lines/s, 26 s for a million lines.connect_read_pipewith< file. Measured:ValueError: Pipe transport is for pipes/sockets only.- Leaving stdin non-blocking. It breaks the next program that reads the terminal.
- Chunk boundaries. Keep the partial last line of each chunk for the next one.
Frequently Asked Questions¶
How do I read stdin asynchronously in Python?
Read chunks in a thread with await asyncio.to_thread(sys.stdin.buffer.read1, 65536) and split them into lines on the loop, keeping the partial last line. It read a million lines in 0.03–0.04 s from both pipes and files.
Why is asyncio.to_thread(sys.stdin.readline) so slow?
Every line pays a thread-pool round trip. Measured at 38,000 lines per second, it took 26 s for a million lines that chunked reading handled in under 0.05 s.
Can I use connect_read_pipe with sys.stdin?
Only when stdin is a pipe, socket or terminal; with a redirected regular file it raised ValueError: Pipe transport is for pipes/sockets only. On a terminal it also left stdin non-blocking after closing, so restore it with os.set_blocking(0, True).
How do I process stdin lines concurrently without running out of memory?
Put lines on a bounded asyncio.Queue consumed by a fixed number of worker tasks; when the queue is full the reader waits, which in turn blocks the upstream program through the pipe.
Related¶
- Async Scripts & CLIs — up to the topic overview.
- Showing progress for concurrent tasks — telling the user how far the pipeline got.
- Asyncio Fundamentals & Event Loop Architecture — the section overview.