StackPractices
intermediate Por Mathias Paulenko

Patrón Pipes and Filters: Pipelines de Datos Componibles

Encadena filtros independientes con pipes para construir pipelines de transformación de datos reutilizables y componibles.

Visión General

El patrón Pipes and Filters toma un trabajo de procesamiento complejo y lo divide en una secuencia de pasos pequeños e independientes. Cada paso es un filtro — datos in, transformación, datos out. Los pipes son solo conectores; en código pueden ser tan simples como composición de funciones o tan complejos como queues y streams. Los filtros terminan siendo reutilizables, componibles y fáciles de testear en aislamiento.

La primera vez que usé este patrón fue arreglando un pipeline de datos en una fintech. Alguien había escrito una función de 300 líneas que parseaba uploads CSV, normalizaba números de teléfono, validaba contra una base de datos de fraude, deduplicaba por email y exportaba a JSON. Cuando la API de fraude empezó a timeout, la función entera crasheaba y perdíamos los datos parseados — todo estaba acoplado. Separarlo en filtros nos permitió reintentar solo el check de fraude sin re-parsear, y testear cada paso sin mockear cinco dependencias. Ese es el valor central: los filtros son chicos, independientes y intercambiables.

Este patrón aparece en todos lados donde los datos se mueven por etapas — jobs ETL, enriquecimiento, cadenas de requests/responses, procesamiento de logs, CI/CD.

Cuándo Usarlo

Pipes and Filters funciona bien cuando una tarea se divide naturalmente en pasos secuenciales e independientes que podés querer reordenar, agregar o remover sin reescribir todo. Es una buena opción cuando la misma transformación aparece en distintos pipelines, o cuando querés testear cada paso por separado. Trabajos ETL, enriquecimiento de datos y cadenas de transformación de requests/responses son usos comunes.

Cuándo NO Usarlo

No optes por este patrón cuando el flujo son dos o tres pasos fijos que nunca cambian: una función simple es más sencilla. Si los pasos están fuertemente acoplados y no se pueden separar en inputs y outputs limpios, Pipes and Filters va a pelear con vos. Y si necesitás que un handler decida si detener el procesamiento, Chain of Responsibility suele ajustarse mejor.

Solución

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);
    }
}

Pipeline async (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())

Explicación

El patrón descompone un trabajo grande en componentes autocontenidos. Un filtro es un paso de procesamiento que recibe datos, los transforma y produce output. Los mejores filtros son funciones puras sin side effects. Los pipes conectan filtros — en un script es composición de funciones; en un sistema distribuido puede ser un topic de Kafka o un pipe de Unix.

Un pipeline es una cadena de filtros conectados por pipes, y como un pipeline se comporta a su vez como un filtro, podés componer pipelines en pipelines más grandes. Esa componibilidad significa que podés reordenar, agregar o remover filtros para armar nuevos pipelines sin tocar los existentes.

Cómo fluyen los datos por un pipeline

flowchart diagram: Input Raw

Trade-offs

El patrón tiene su costo. Cada filtro agrega una llamada a función y potencialmente una copia de datos. Para sistemas de alto throughput que procesan millones de registros, ese overhead se acumula. En Python 3.12+, los generadores e itertools pueden ayudar haciendo que los filtros sean lazy — procesan un item a la vez en vez de materializar la lista completa. En Java 21+, Stream ya te da evaluación lazy gratis. En Node 20+, podés usar Node streams o async generators para pipelines con backpressure.

Algo que me mordió: el manejo de errores se vuelve más difícil, no más fácil. Cuando un filtro falla a mitad del pipeline, tenés tres opciones — saltar el item, reintentar o abortar el pipeline entero. Envolver el pipeline en un try/catch único funciona para casos simples, pero los sistemas de producción suelen necesitar boundaries de error por filtro con dead-letter queues para los items que fallan repetidamente.

Variantes

VarianteEjecuciónCaso de uso
Pipeline sincrónicoSecuencial, bloqueanteTransformación simple de datos
Pipeline asyncNo bloqueante, concurrenteProcesamiento I/O-bound (HTTP, DB)
Pipeline paraleloFiltros corren en paraleloTransformaciones CPU-bound
Pipeline streamingEvent-driven, continuoStreams de datos en tiempo real
Pipeline batchProcesa en chunksETL, procesamiento de datos programado

Para streaming en tiempo real, cuidado con filtros lentos. El patrón Back Pressure muestra cómo evitar que productores rápidos saturen consumidores lentos.

En la práctica, vas a encontrar este patrón en herramientas como Apache Airflow (DAGs de filtros), Kafka Streams (pipelines de stream processing) e incluso en pipes de shell — cat file | grep | sort | uniq es básicamente un pipeline de Pipes and Filters a nivel OS. Para un análisis más profundo de pipelines de streaming, consultá la Guía de Kafka Stream Processing.

Buenas Prácticas

  • Mantené los filtros puros, sin side effects ni estado mutable compartido, así siguen siendo testeables y componibles.
  • Un filtro, un trabajo. No metas múltiples transformaciones en un solo filtro — eso los hace más fáciles de entender y reusar.
  • Usá firmas de tipos para documentar el contrato entre filtros.
  • Manejá errores a nivel del pipeline en vez de dentro de cada filtro, así no ocultás fallas.
  • Construí pipelines dinámicamente con un builder o configuración cuando el orden no es fijo.
  • Testeá los filtros en aislamiento; las funciones puras son fáciles de testear unitariamente.
  • Insertá filtros de logging entre pasos de procesamiento para debugging sin tocar la lógica de negocio.

Errores Comunes

  • Hacer filtros stateful. El estado compartido rompe la composabilidad y dificulta los tests.
  • Meter side effects dentro de los filtros. Escribir en una base de datos o llamar una API desde un filtro hace el pipeline no determinista.
  • Ignorar errores. Un fallo no manejado en un filtro puede romper todo el pipeline.
  • Fijar el orden de filtros en el código. Usá un builder o configuración para que el pipeline pueda evolucionar.
  • Filtros que intentan hacer de todo. Un filtro tiene que hacer una sola transformación, no más.
  • No tipar inputs y outputs de filtros. Los errores de tipo en runtime son un dolor de cabeza para debuggear.
  • Olvidar la backpressure en pipelines streaming. Los filtros lentos hacen que la memoria se acumule en los pipes hasta que el proceso se queda sin RAM.

Resumen

  • Los filtros son funciones puras — sin side effects ni estado compartido. Eso los hace testeables, componibles y seguros de reordenar.
  • Los pipes son conectores — composición de funciones en casos simples, queues o streams en sistemas distribuidos.
  • La componibilidad es la gran ventaja — los pipelines se comportan como filtros, así que podés anidarlos y reusarlos.
  • El manejo de errores va a nivel del pipeline, no dentro de cada filtro. Envolver el pipeline en try/catch o devolver Result<T, E> — así el caller elige la estrategia de recuperación.
  • Cuidado con el overhead — cada filtro agrega un boundary de llamada. Usá evaluación lazy (generators, streams) para pipelines de alto throughput.

Código complementario: El repo companion de pipes-and-filters-pattern tiene ejemplos ejecutables en Python, JavaScript y Java más tests.

See Also

Preguntas frecuentes

¿En qué se diferencia de Chain of Responsibility?

En Chain of Responsibility, cada handler decide si pasa el request o se detiene. En Pipes and Filters, cada filtro procesa los datos y los pasa al siguiente. Pipes and Filters es sobre transformación; Chain of Responsibility es sobre manejo. Consultá Chain of Responsibility Pattern para la diferencia.

¿Uso esto o una función simple?

Usá Pipes and Filters cuando el orden puede cambiar, cuando los mismos filtros aparecen en varios pipelines, o cuando querés testear cada paso por separado. Si la tarea son dos o tres pasos fijos que nunca cambian, una función simple suele ser suficiente.

¿Cómo manejo branching en un pipeline?

Usá un router filter que envíe datos a distintos sub-pipelines según una condición. El router es a su vez un filtro: recibe input, evalúa una condición y rutea al sub-pipeline adecuado.

¿Cómo manejo errores en un pipeline?

Envolvé el pipeline completo en un try/catch o usá un tipo resultado, como Result<T, E>. Dejá que el caller decida cómo manejar un paso fallido. No captures errores dentro de filtros individuales — si lo hacés, las fallas se ocultan silenciosamente.

¿Cuándo uso un pipeline async o paralelo?

Usá un pipeline async cuando los filtros esperan I/O. Usá uno paralelo cuando los filtros son CPU-bound y pueden correr independientemente. Para streams en tiempo real, usá un pipeline streaming con manejo de backpressure.