Sooner or later a CLI needs the equivalent of git log --format=%ae | sort | uniq -c, pg_dump mydb | gzip, or tar -c dir | ssh host tar -x. The tempting shortcut is subprocess.run("cmd1 | cmd2", shell=True). It works until a file name contains a space or a quote, a user-supplied value turns into shell syntax, or the first command fails and the shell reports success because only the last command’s status counts. Building the pipeline in Python — one Popen per stage, connected by OS pipes — avoids the shell entirely, keeps arguments as lists, and lets you check every stage’s exit status. The details are fiddly: which pipe ends to close, how to avoid deadlocks on unread stderr, what to do when a downstream command exits early, and how to clean up on timeouts. This guide writes a small run_pipeline helper that handles them, and tests each failure mode. It belongs to the running subprocesses from Python CLIs topic.
Prerequisites
- Calling external commands safely with subprocess and avoiding shell injection in Python CLIs.
- A Unix-like system for the examples; the helper itself also runs on Windows with Windows commands.
How a pipeline is wired
A shell pipeline is nothing magic: the shell creates an OS pipe, starts the first command with its stdout connected to the pipe’s write end, starts the second with its stdin connected to the read end, and waits for both. Data streams between the processes directly, at full speed, without passing through the shell. Python can do the same: Popen(..., stdout=PIPE) creates the pipe, and passing that process’s .stdout as the next process’s stdin connects them. Your own process only reads the final stage’s output.
The recipe
# src/mytool/pipeline.py
from __future__ import annotations
import subprocess
import tempfile
from collections.abc import Sequence
from dataclasses import dataclass
@dataclass
class PipelineError(Exception):
command: Sequence[str]
returncode: int
stderr: str
def __str__(self) -> str:
detail = f": {self.stderr.strip().splitlines()[-1]}" if self.stderr.strip() else ""
return f"{self.command[0]} exited with status {self.returncode}{detail}"
def run_pipeline(*commands: Sequence[str], input: bytes | None = None, timeout: float | None = None,
ok_codes: dict[int, set[int]] | None = None) -> bytes:
"""Run cmd1 | cmd2 | ... without a shell and return the last command's stdout.
Every command's exit status is checked (like `set -o pipefail`); ok_codes lists extra
acceptable codes per position, e.g. {1: {1}} for a grep that may match nothing.
"""
ok_codes = ok_codes or {}
procs: list[subprocess.Popen] = []
# stdin and every stderr go through temporary files: pipes we are not reading can fill
# up and block a stage forever, files cannot.
stdin = tempfile.TemporaryFile() if input is not None else subprocess.DEVNULL
stderrs = [tempfile.TemporaryFile() for _ in commands]
try:
if input is not None:
stdin.write(input)
stdin.seek(0)
for i, command in enumerate(commands):
proc = subprocess.Popen(command, stdin=procs[-1].stdout if procs else stdin,
stdout=subprocess.PIPE, stderr=stderrs[i])
if procs:
procs[-1].stdout.close() # the next stage now holds the only read end
procs.append(proc)
out, _ = procs[-1].communicate(timeout=timeout)
for proc in procs[:-1]:
proc.wait(timeout=timeout)
except BaseException: # timeout, Ctrl+C, failure to start a stage
for proc in procs:
proc.kill()
proc.wait()
raise
finally:
if input is not None:
stdin.close()
try:
for i, proc in enumerate(procs):
if proc.returncode != 0 and proc.returncode not in ok_codes.get(i, set()):
stderrs[i].seek(0)
raise PipelineError(commands[i], proc.returncode,
stderrs[i].read().decode(errors="replace"))
finally:
for f in stderrs:
f.close()
return out
Used from a command:
from mytool.pipeline import run_pipeline
authors = run_pipeline(
["git", "log", "--format=%ae"],
["sort"],
["uniq", "-c"],
timeout=60,
)
for line in authors.decode().splitlines():
print(line.strip())
Closing the parent’s copy of each pipe
After starting stage n+1 with stage n’s stdout as its stdin, the helper closes its own handle to that pipe. This line is the one most hand-written pipelines miss, and it matters in both directions. A pipe’s reader sees end-of-file only when every write end is closed, and a writer gets SIGPIPE only when every read end is closed. If your Python process keeps a copy of the read end open, a downstream command that exits early (head -n1) does not cause the upstream command to stop — it keeps producing into a pipe nobody reads, fills it, and blocks forever. Closing the parent’s copy leaves the next stage as the only reader, which is how the shell does it.
Checking every stage, like pipefail
A shell pipeline’s exit status is the last command’s, so false | sort succeeds. Bash’s set -o pipefail changes that; the helper does the equivalent by checking every process’s return code and raising PipelineError for the first failing stage, including the last line of its stderr. Some non-zero codes are not failures: grep exits 1 when nothing matches, diff exits 1 when files differ. ok_codes lets callers list acceptable codes per stage position rather than disabling checks for the whole pipeline.
Early exits and SIGPIPE
When a downstream command quits before reading everything, the upstream one is killed by SIGPIPE the next time it writes: yes | head -n1 ends with yes reporting status −13 in Python’s terms (killed by signal 13). That is normal behaviour, not an error, for producers whose output is deliberately cut short. Popen restores the default SIGPIPE handling in children (restore_signals=True), so ordinary Unix tools behave exactly as they do in a shell. Allow −13 for producer stages you expect to be cut off, and see handling broken pipe and SIGPIPE for the same issue in your own output.
Avoiding deadlocks
A process blocks when it writes into a full pipe whose reader is not reading. Python only reads the final stage’s stdout, so any other pipe it holds — an intermediate stage’s stderr, or the first stage’s stdin — can fill up and block the whole chain. The helper sidesteps this by using temporary files for every stage’s stderr and for the input data: files never fill, so no stage can block on them, and the stderr content is still available for the error message. communicate() on the last stage reads its stdout to the end, after which the earlier stages have finished or are about to.
Timeouts and cleanup
communicate(timeout=...) raises TimeoutExpired when the pipeline takes too long. The except BaseException block then kills every stage and reaps it, so no orphaned processes keep running in the background — the same applies to Ctrl+C, which arrives as KeyboardInterrupt. The approach is the one described in handling subprocess timeouts and exit codes, extended to several processes.
When to use Python instead of a stage
Not every stage needs to be a process. sort | uniq -c is a collections.Counter; gzip is the gzip module; grep pattern is a list comprehension. Replacing stages with Python code removes dependencies on external tools (and their platform differences — BSD versus GNU sort, missing tools on Windows) and makes errors easier to report. Keep external stages for what Python cannot do as well: pg_dump, ssh, tar with special file handling, compressors like zstd, or tools the user expects to see invoked.
UX considerations
- Report which stage failed. “sort exited with status 2: sort: cannot read: …” beats “pipeline failed”.
- Stream large outputs. The helper collects the final output in memory, which is fine for reports; for gigabytes, read
procs[-1].stdoutincrementally or connect the last stage directly to a file opened by your CLI. - Show the equivalent shell command at
-v.shlex.joineach stage and join with|, so users can reproduce problems by hand. - Check that tools exist first.
shutil.which("pg_dump")before starting gives a clear “pg_dump is not installed” instead of aFileNotFoundErrorhalfway through.
Testing the behaviour
Use the Python interpreter itself to simulate stages, so the tests run anywhere and exercise exactly the failure modes that matter — upstream failures, allowed codes, early exits, noisy stderr and timeouts:
# tests/test_pipeline.py
import shutil
import subprocess
import sys
import pytest
from mytool.pipeline import PipelineError, run_pipeline
PY = sys.executable
pytestmark = pytest.mark.skipif(not shutil.which("sort"), reason="needs sort")
def test_two_stage_pipeline_with_input():
out = run_pipeline([PY, "-c", "import sys; sys.stdout.write(sys.stdin.read().upper())"],
["sort"], input=b"b\na\nc\n")
assert out == b"A\nB\nC\n"
def test_failure_in_first_stage_is_not_hidden():
with pytest.raises(PipelineError) as info:
run_pipeline([PY, "-c", "import sys; print('partial'); sys.exit('boom')"], ["sort"])
assert info.value.returncode == 1 and "boom" in str(info.value)
def test_allowed_exit_codes():
out = run_pipeline([PY, "-c", "print('x')"],
[PY, "-c", "import sys; sys.stdin.read(); sys.exit(1)"],
ok_codes={1: {1}})
assert out == b""
def test_early_exit_downstream_does_not_hang():
# The consumer quits after one line; the producer then fails writing (EPIPE), which we accept.
producer = [PY, "-c", "import sys\nfor i in range(10**6): print(i)"]
consumer = [PY, "-c", "import sys; print(sys.stdin.readline().strip())"]
out = run_pipeline(producer, consumer, ok_codes={0: {1, -13}}, timeout=20)
assert out == b"0\n"
def test_timeout_kills_every_stage():
with pytest.raises(subprocess.TimeoutExpired):
run_pipeline([PY, "-c", "import time; time.sleep(30)"], ["sort"], timeout=0.5)
def test_noisy_stderr_does_not_deadlock():
noisy = [PY, "-c", "import sys\nfor i in range(20000): print('warn', i, file=sys.stderr)\nprint('ok')"]
assert run_pipeline(noisy, ["sort"], timeout=20) == b"ok\n"
The early-exit test would hang without the “close the parent’s copy” line, which is why it has a timeout; the noisy-stderr test would hang without temporary files for stderr. Both turn subtle deadlocks into fast, deterministic test failures.
Conclusion
Pipelines do not need a shell. Start one Popen per stage, pass each stage’s stdout as the next stage’s stdin and close the parent’s copy, send stderr and input through temporary files so nothing can deadlock, read the final output with communicate(timeout=...), check every stage’s exit status with per-stage exceptions for codes like grep’s 1 and SIGPIPE’s −13, kill and reap every stage on timeouts and Ctrl+C, and replace stages with Python where that is simpler.
Frequently asked questions
Is shell=True ever acceptable for pipelines?
For a fixed command string with no user-supplied parts, run in a context you control, it is safe — but you still lose per-stage exit codes unless you add set -o pipefail, and it requires a POSIX shell. The explicit form is safer by construction and portable.
How do I write the pipeline’s output straight to a file?
Pass an open file as the last stage’s stdout instead of PIPE: open("dump.sql.gz", "wb"). Then communicate() returns no data, and the file receives the output without passing through Python.
Can I feed a large input without temporary files?
Yes, by writing to the first stage’s stdin from a separate thread while the main thread reads the last stage’s stdout — which is what communicate() does internally for a single process. A temporary file is simpler and fast enough for most CLIs.
Does this work with asyncio?
asyncio.create_subprocess_exec accepts the same idea, but connecting one process’s output to another’s input needs os.pipe() file descriptors passed explicitly. For most CLIs, running the synchronous helper in a thread is simpler.
What about Windows?
Popen pipes work the same way on Windows; the commands differ. There is no SIGPIPE, so an upstream process usually gets a write error instead and exits with its own error code.