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
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
| Variante | Ejecución | Caso de uso |
|---|---|---|
| Pipeline sincrónico | Secuencial, bloqueante | Transformación simple de datos |
| Pipeline async | No bloqueante, concurrente | Procesamiento I/O-bound (HTTP, DB) |
| Pipeline paralelo | Filtros corren en paralelo | Transformaciones CPU-bound |
| Pipeline streaming | Event-driven, continuo | Streams de datos en tiempo real |
| Pipeline batch | Procesa en chunks | ETL, 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
- Enterprise Integration Patterns — Pipes and Filters — Referencia canónica de Hohpe & Woolf para el patrón.
- Martin Fowler: Collection Pipeline — Cómo los lenguajes modernos componen filtros sobre colecciones.
- Documentación de Python asyncio — async/await para pipelines no bloqueantes en Python 3.12+.
- Java Stream API — Primitivas de pipeline integradas en Java 21+.
- Apache Airflow — Orquestación de pipelines de producción con DAGs.
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.
Recursos Relacionados
Patrón Chain of Responsibility
Pasa solicitudes a lo largo de una cadena de manejadores hasta que uno la procese. Un patrón de comportamiento para desacoplar emisores y receptores.
PatternPatrón Observer
Define un mecanismo de suscripción para notificar a múltiples objetos sobre eventos. Patrón de diseño conductual para comunicación basada en eventos.
PatternPatrón Back-Pressure
Previene que sistemas upstream abrumen a consumidores downstream propagando señales de control de flujo hacia atrás a través del pipeline, asegurando throughput estable bajo carga.
PatternPatrón Marker Interface
Usa interfaces vacías como tags de metadata para señalar propiedades o capacidades en tiempo de compilación y runtime, habilitando verificaciones type-safe sin modificar comportamiento de clase.
GuideKafka Stream Processing
Construye pipelines de event streaming en tiempo real con Kafka. Cubre producers, consumers, Kafka Streams, Kafka Connect, schema registry y patrones de procesamiento.
PatternPatrón de Hosting de Contenido Estático
Despliega archivos estáticos en una red de entrega de contenido o almacenamiento de objetos para descargar servidores de origen, reducir latencia y mejorar disponibilidad.