StackPractices
intermediate By Mathias Paulenko

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

flowchart diagram: Raw Input

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

VariantExecutionUse Case
Synchronous pipelineSequential, blockingSimple data transformation
Async pipelineNon-blocking, concurrentI/O-bound processing (HTTP, DB)
Parallel pipelineFilters run in parallelCPU-bound transformations
Streaming pipelineEvent-driven, continuousReal-time data streams
Batch pipelineProcess in chunksETL, 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

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.