Pipes and Filters Pattern: Composable Data Pipelines
Chain processing steps with independent filters connected by pipes. A pattern for data transformation pipelines where each step is reusable and composable.
Overview
The Pipes and Filters Pattern takes a complex processing job and breaks it into small, independent steps. Each step is a filter — data in, transform, data out. Pipes are just connectors; in code they can be as simple as function composition or as complex as queues and streams. Filters end up reusable, composable, and easy to test in isolation.
I first ran into this pattern debugging a data pipeline at a fintech startup. Someone had written a 300-line function that parsed CSV uploads, normalized phone numbers, validated against a fraud database, deduplicated by email, and exported to JSON. When the fraud API started timing out, the whole function crashed and we lost the parsed data — everything was coupled together. Splitting it into filters meant we could retry just the fraud check without re-parsing, and we could unit test each step without mocking five dependencies. That’s the core value: filters are small, independent, and swappable.
You’ll see this pattern wherever data moves through stages — ETL jobs, enrichment, request/response chains, log processing, CI/CD.
When to Use
Pipes and Filters works well when a task naturally splits into sequential, independent steps that you might want to reorder, add, or remove without rewriting everything. It’s a good choice when the same transformation shows up in different pipelines, or when you want to unit test each step on its own. ETL jobs, data enrichment, and request/response transformation chains are common fits.
When NOT to Use
Don’t reach for this pattern when the workflow is just two or three fixed steps that never change — a plain function is simpler. If the steps are tightly coupled and can’t be separated into clean inputs and outputs, Pipes and Filters will fight you. And if you need a handler to decide whether to stop processing, Chain of Responsibility is usually a better fit.
Solution
Python
from typing import Callable, Any
from dataclasses import dataclass
Filter = Callable[[Any], Any]
def pipe(*filters: Filter) -> Filter:
def pipeline(data: Any) -> Any:
result = data
for f in filters:
result = f(result)
return result
return pipeline
# Filters — each is a pure function
def parse_csv(raw: str) -> list[dict]:
lines = raw.strip().split("\n")
headers = lines[0].split(",")
return [
dict(zip(headers, line.split(",")))
for line in lines[1:]
]
def filter_active(records: list[dict]) -> list[dict]:
return [r for r in records if r.get("status") == "active"]
def normalize_emails(records: list[dict]) -> list[dict]:
for r in records:
r["email"] = r.get("email", "").lower().strip()
return records
def deduplicate(records: list[dict]) -> list[dict]:
seen = set()
result = []
for r in records:
key = r.get("email")
if key not in seen:
seen.add(key)
result.append(r)
return result
def to_json(records: list[dict]) -> str:
import json
return json.dumps(records, indent=2)
# Compose a pipeline
process_users = pipe(
parse_csv,
filter_active,
normalize_emails,
deduplicate,
to_json,
)
# Usage
raw_data = """name,email,status
Alice,ALICE@Example.COM,active
Bob,bob@example.com,inactive
Charlie,CHARLIE@example.com,active
Alice,alice@example.com,active"""
result = process_users(raw_data)
print(result)
JavaScript
function pipe(...filters) {
return (data) => filters.reduce((acc, fn) => fn(acc), data);
}
// Filters — each is a pure function
function parseCsv(raw) {
const lines = raw.trim().split("\n");
const headers = lines[0].split(",");
return lines.slice(1).map((line) => {
const values = line.split(",");
return Object.fromEntries(headers.map((h, i) => [h, values[i]]));
});
}
function filterActive(records) {
return records.filter((r) => r.status === "active");
}
function normalizeEmails(records) {
return records.map((r) => ({
...r,
email: (r.email || "").toLowerCase().trim(),
}));
}
function deduplicate(records) {
const seen = new Set();
return records.filter((r) => {
if (seen.has(r.email)) return false;
seen.add(r.email);
return true;
});
}
function toJson(records) {
return JSON.stringify(records, null, 2);
}
// Compose a pipeline
const processUsers = pipe(
parseCsv,
filterActive,
normalizeEmails,
deduplicate,
toJson
);
// Usage
const rawData = `name,email,status
Alice,ALICE@Example.COM,active
Bob,bob@example.com,inactive
Charlie,CHARLIE@example.com,active
Alice,alice@example.com,active`;
console.log(processUsers(rawData));
Java
import java.util.*;
import java.util.function.Function;
import java.util.stream.Collectors;
public class PipesAndFilters {
@FunctionalInterface
interface Filter<T, R> extends Function<T, R> {}
static <T> Filter<T, T> pipe(Filter<T, T>... filters) {
return data -> {
T result = data;
for (Filter<T, T> f : filters) {
result = f.apply(result);
}
return result;
};
}
// Filters
static Filter<String, List<Map<String, String>>> parseCsv = raw -> {
String[] lines = raw.trim().split("\n");
String[] headers = lines[0].split(",");
return Arrays.stream(lines, 1, lines.length)
.map(line -> {
String[] values = line.split(",");
Map<String, String> record = new LinkedHashMap<>();
for (int i = 0; i < headers.length; i++) {
record.put(headers[i], values[i]);
}
return record;
})
.collect(Collectors.toList());
};
static Filter<List<Map<String, String>>, List<Map<String, String>>> filterActive =
records -> records.stream()
.filter(r -> "active".equals(r.get("status")))
.collect(Collectors.toList());
static Filter<List<Map<String, String>>, List<Map<String, String>>> normalizeEmails =
records -> records.stream()
.map(r -> {
r.put("email", r.get("email").toLowerCase().trim());
return r;
})
.collect(Collectors.toList());
static Filter<List<Map<String, String>>, List<Map<String, String>>> deduplicate =
records -> {
Set<String> seen = new HashSet<>();
return records.stream()
.filter(r -> seen.add(r.get("email")))
.collect(Collectors.toList());
};
public static void main(String[] args) {
String rawData = "name,email,status\n" +
"Alice,ALICE@Example.COM,active\n" +
"Bob,bob@example.com,inactive\n" +
"Charlie,CHARLIE@example.com,active";
var pipeline = pipe(parseCsv, filterActive, normalizeEmails, deduplicate);
List<Map<String, String>> result = pipeline.apply(rawData);
result.forEach(System.out::println);
}
}
Async pipeline (Python)
import asyncio
from typing import Any, Callable, Awaitable
AsyncFilter = Callable[[Any], Awaitable[Any]]
async def async_pipe(*filters: AsyncFilter) -> AsyncFilter:
async def pipeline(data: Any) -> Any:
result = data
for f in filters:
result = await f(result)
return result
return pipeline
async def fetch_data(url: str) -> dict:
await asyncio.sleep(0.1) # simulate HTTP
return {"url": url, "status": 200, "body": "raw data"}
async def parse_data(raw: dict) -> dict:
await asyncio.sleep(0.05)
raw["parsed"] = raw["body"].upper()
return raw
async def validate_data(data: dict) -> dict:
await asyncio.sleep(0.05)
if data["status"] != 200:
raise ValueError(f"Bad status: {data['status']}")
data["valid"] = True
return data
async def enrich_data(data: dict) -> dict:
await asyncio.sleep(0.05)
data["enriched"] = f"ENRICHED:{data['parsed']}"
return data
async def main():
pipeline = await async_pipe(fetch_data, parse_data, validate_data, enrich_data)
result = await pipeline("https://api.example.com/data")
print(result)
asyncio.run(main())
Explanation
The pattern decomposes a big job into self-contained components. Think of a filter as one step: data in, transform, data out. The best filters are pure functions with no side effects. Pipes connect filters — in a script that’s just function composition; in a distributed system it might be a Kafka topic or a Unix pipe.
A pipeline is a chain of filters connected by pipes, and because a pipeline itself behaves like a filter, you can compose pipelines into larger pipelines. That composability means you can reorder, add, or remove filters to build new pipelines without touching existing ones.
How data flows through a pipeline
Trade-offs
The pattern isn’t free. Each filter adds a function call and potentially a data copy. If you’re
processing millions of records, that per-filter overhead adds up fast. In Python 3.12+,
generators and itertools can help by making filters lazy — they process one item at a time
instead of materializing the full list. In Java 21+, Stream already gives you lazy evaluation
for free. In Node 20+, you can use Node streams or async generators for backpressure-aware
pipelines.
One thing that bit me: error handling gets harder, not easier. When a filter blows up mid-pipeline, you’ve got three choices — skip the item, retry it, or kill the whole pipeline. Wrapping everything in a single try/catch works for simple scripts, but production systems usually need per-filter error boundaries with dead-letter queues for items that fail repeatedly.
Variants
| Variant | Execution | Use Case |
|---|---|---|
| Synchronous pipeline | Sequential, blocking | Simple data transformation |
| Async pipeline | Non-blocking, concurrent | I/O-bound processing (HTTP, DB) |
| Parallel pipeline | Filters run in parallel | CPU-bound transformations |
| Streaming pipeline | Event-driven, continuous | Real-time data streams |
| Batch pipeline | Process in chunks | ETL, scheduled data processing |
For real-time streaming, watch out for slow filters. The Back Pressure Pattern shows how to keep fast producers from overwhelming slow consumers.
In practice, you’ll find this pattern in tools like Apache Airflow (DAGs of filters), Kafka Streams
(stream processing pipelines), and even plain shell pipes — cat file | grep | sort | uniq is a
Pipes and Filters pipeline at the OS level. For a deeper dive on streaming pipelines, see the
Kafka Stream Processing Guide.
Best Practices
- Keep filters pure — no side effects and no shared mutable state — so they stay testable and composable.
- One filter, one job. Don’t cram multiple transformations into a single filter — that’s the sweet spot for readability and reuse.
- Use type signatures to document the contract between filters.
- Don’t catch errors inside individual filters — if you do, failures get swallowed silently.
- Build pipelines dynamically with a builder or configuration when the filter order isn’t fixed.
- Test filters in isolation; pure functions are straightforward to unit test.
- Slip a logging filter between steps when you need to debug — no need to touch business logic.
Common Mistakes
- Making filters stateful. Shared state breaks composability and makes tests harder.
- Putting side effects inside filters. Writing to a database or calling an API from a filter makes the pipeline non-deterministic.
- Ignoring errors. One unhandled failure in a filter can crash the whole pipeline.
- Hardcoding filter order. Use a builder or configuration so the pipeline can evolve.
- Overloading a filter. When a filter tries to do several transformations, it becomes harder to test and reuse.
- Skipping types on filter inputs and outputs. Runtime type mismatches become painful to debug.
- Forgetting backpressure in streaming pipelines. Slow filters cause memory to pile up in the pipes until the process runs out of RAM. Slow filters can cause memory to build up in pipes.
Summary
- Filters are pure functions — no side effects, no shared state. That makes them testable, composable, and safe to reorder.
- Pipes are connectors — function composition in simple cases, queues or streams in distributed systems.
- Composability is the big win — pipelines behave like filters, so you can nest and reuse them.
- Handle errors at the pipeline level, not inside each filter. Wrap the whole pipeline in
try/catch or return
Result<T, E>— that way the caller picks the recovery strategy. - Watch for overhead — each filter adds a call boundary. Use lazy evaluation (generators, streams) for high-throughput pipelines.
Companion code: The pipes-and-filters-pattern companion repo has runnable Python, JavaScript, and Java examples plus tests.
See Also
- Enterprise Integration Patterns — Pipes and Filters — Hohpe & Woolf’s canonical reference for the pattern.
- Martin Fowler: Collection Pipeline — how modern languages compose filters over collections.
- Python asyncio documentation — async/await for non-blocking pipelines in Python 3.12+.
- Java Stream API — built-in pipeline primitives in Java 21+.
- Apache Airflow — production-grade pipeline orchestration with DAGs.
Frequently Asked Questions
How is this different from Chain of Responsibility?
In Chain of Responsibility, each handler decides whether to pass the request along or stop. In Pipes and Filters, every filter processes the data and passes it to the next. Pipes and Filters is about transformation; Chain of Responsibility is about handling. See Chain of Responsibility Pattern for the distinction.
Should I use this or a simple function?
Reach for Pipes and Filters when the order might change, when the same filters pop up in different pipelines, or when you want to test each step independently. If the task is two or three fixed steps that never change, a plain function is usually enough.
How do I handle branching in a pipeline?
Drop in a router filter — it checks a condition and sends data to the right sub-pipeline. The router is itself a filter: it receives input, evaluates a condition, and routes accordingly.
How do I handle errors in a pipeline?
Wrap the whole pipeline in a try/catch or use a result type such as Result<T, E>. Let the caller
decide
how to handle a failed step. Avoid catching errors inside individual filters, or you'll hide
failures.
When should I use an async or parallel pipeline?
Use an async pipeline when filters wait on I/O. Use a parallel pipeline when filters are CPU-bound and can run independently. Got real-time streams? Go with a streaming pipeline, and don't forget backpressure.
Related Resources
Chain of Responsibility Pattern
Pass requests along a chain of handlers until one handles it. A behavioral design pattern for decoupling senders and receivers.
PatternObserver Pattern
Define a subscription mechanism to notify multiple objects about events. A behavioral design pattern for event-driven communication.
PatternBack-Pressure Pattern
Prevent upstream systems from overwhelming downstream consumers by propagating flow-control signals backward through the pipeline, ensuring stable throughput under load.
PatternMarker Interface Pattern
Use empty interfaces as metadata tags to signal properties or capabilities at compile time and runtime, enabling type-safe checks without modifying class behavior.
GuideComplete Guide to Kafka Stream Processing
Build real-time event streaming pipelines with Kafka. Covers producers, consumers, Kafka Streams, Kafka Connect, schema registry, and stream processing patterns.
PatternStatic Content Hosting Pattern
Deploy static files to a dedicated content delivery network or object storage to offload origin servers, reduce latency, and improve availability for assets like images, CSS, and JavaScript.