Skip to content

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

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.

Reading 1,000,000 lines from a pipe 4 horizontal bars comparing to_thread(readline) per line with the others. Reading 1,000,000 lines from a pipe to_thread(readline) per line 38k lines/s connect_read_pipe + StreamReader 1.69M lines/s reader thread -> Queue, batches of 1,000 14.4M lines/s to_thread(read1(64 KiB)) chunks 28.2M lines/s Values in millions of lines per second; 25.8 MB file, Python 3.14. Chunk and queue rates exclude per-line processing. The thread hop per line, not the I/O, set the slowest rate.

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.

Which reader worked with which stdin A grid of 4 rows by 4 columns. Which reader worked with which stdin reader pipe redirected file terminal to_thread(readline) per line works, slow works, slow works to_thread(read1) chunks works works works connect_read_pipe works ValueError works, leaves O_NONBLOCK set reader thread -> Queue works works works Only the thread-based readers handled every way a script gets run.

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.

A pipeline script with backpressure end to end A flow of 5 stages. A pipeline script with backpressure end to end upstream program writes to the pipe read1(64 KiB) in a thread one hop per chunk split lines on the loop keep partial tail bounded Queue put() waits when full N worker tasks process and record status A full queue stops the reader, and a stopped reader stops the producer.

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_pipe was 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_pipe with < 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.