Data Engineering Design Patterns: 8 patrones de ingesta que tenés que conocer
En este episodio repaso los capítulos 1 y 2 del libro Data Engineering Design Patterns de Bartosz Konieczny, cubriendo los fundamentos de patrones de diseño aplicados a data engineering y los 8 patrones de ingesta de datos que el autor presenta.
El libro es una lectura que recomiendo mucho si ya tenés experiencia construyendo pipelines y querés poner nombre a las cosas que venís haciendo (o que deberías estar haciendo). No es un libro introductorio: asume que sabés qué es un pipeline y que ya te peleaste con problemas reales.
Por qué importan los patrones
Cuando trabajás en un equipo de datos, uno de los problemas más frecuentes no es técnico sino comunicativo. Alguien dice “cargamos todo de nuevo” y otro entiende algo distinto. Los patrones de diseño resuelven eso: te dan un vocabulario compartido para discutir soluciones.
No se trata de memorizar nombres. Se trata de que cuando alguien diga “usemos CDC para esta fuente” o “necesitamos un Compactor después del streaming”, todos en la mesa entiendan exactamente lo mismo. Es la misma lógica que los patrones de diseño en software (Factory, Observer, etc.), pero aplicada a flujos de datos.
Konieczny organiza los patrones en categorías. En este post me enfoco en los 8 patrones de ingesta, que son los que más impacto tienen en el día a día.
Los 8 patrones de ingesta
1. Full Loader
Carga completa de la fuente en cada ejecución. Borrás (o sobreescribís) todo lo que había y traés el dataset entero de nuevo. Es el patrón más simple y el que menos puede fallar, pero el más costoso en tiempo y recursos cuando la tabla crece.
Cuándo usarlo: tablas de referencia chicas (países, monedas, configuraciones), fuentes que no tienen columna de fecha de modificación, o cuando necesitás garantizar consistencia total sin complicarte.
# Full Loader en PySpark — sobreescribir tabla de dimensiones
df_source = (spark.read.format("jdbc")
.option("url", jdbc_url)
.option("dbtable", "erp.dim_currency")
.option("user", user)
.option("password", password)
.load())
(df_source.write
.format("delta")
.mode("overwrite")
.option("overwriteSchema", "true")
.saveAsTable("bronze.dim_currency"))2. Incremental Loader
Solo trae los registros nuevos o modificados desde la última ejecución. Requiere una columna de control confiable: un timestamp de modificación, un ID autoincremental, o algo equivalente. Es el patrón más común en producción porque equilibra eficiencia con simplicidad.
Cuándo usarlo: tablas transaccionales que crecen constantemente y tienen una columna updated_at o similar. Es el caballo de batalla de la mayoría de los pipelines batch.
# Incremental Loader — traer solo registros nuevos
from pyspark.sql import functions as F
# Obtener la última marca de agua (high watermark)
max_timestamp = (spark.table("bronze.orders")
.agg(F.max("updated_at"))
.collect()[0][0])
# Leer solo lo nuevo desde la fuente
df_incremental = (spark.read.format("jdbc")
.option("url", jdbc_url)
.option("dbtable", "erp.orders")
.option("user", user)
.option("password", password)
.load()
.filter(F.col("updated_at") > max_timestamp))
# Append o merge en destino
(df_incremental.write
.format("delta")
.mode("append")
.saveAsTable("bronze.orders"))3. CDC (Change Data Capture)
Captura cambios directamente desde el log de transacciones de la base de datos origen. Es el patrón más eficiente para fuentes transaccionales porque no necesitás escanear la tabla origen: leés los cambios del WAL (Write-Ahead Log) o el equivalente del motor. Herramientas como Debezium, Fivetran o el conector CDC de Databricks hacen el trabajo pesado.
Cuándo usarlo: cuando la fuente es una base de datos relacional con alto volumen de cambios, cuando necesitás capturar deletes (que el Incremental Loader no ve), o cuando querés latencia baja sin impactar la fuente.
-- CDC con MERGE en Delta Lake
-- Los eventos CDC llegan con una columna _operation: INSERT, UPDATE, DELETE
MERGE INTO silver.customers AS target
USING (
SELECT * FROM bronze.customers_cdc
QUALIFY ROW_NUMBER() OVER (
PARTITION BY customer_id ORDER BY _event_timestamp DESC
) = 1
) AS source
ON target.customer_id = source.customer_id
WHEN MATCHED AND source._operation = 'DELETE' THEN DELETE
WHEN MATCHED AND source._operation = 'UPDATE' THEN UPDATE SET *
WHEN NOT MATCHED AND source._operation != 'DELETE' THEN INSERT *;4. Passthrough Replicator
Replica datos tal cual vienen de la fuente, sin ninguna transformación. Es como un espejo: lo que hay en origen, lo tenés igual en destino. Suena trivial, pero es un patrón deliberado: separar la extracción de la transformación te da un punto de recuperación limpio.
Cuándo usarlo: capas de staging o Bronze donde querés mantener una copia fiel de la fuente. Si algo sale mal en la transformación posterior, siempre podés volver a los datos originales sin pegar de nuevo al sistema origen.
5. Transformation Replicator
Replica y transforma en el mismo paso. A diferencia del Passthrough, acá aplicás lógica de negocio, limpieza o reestructuración durante la ingesta misma. Reduce la cantidad de jobs pero acopla extracción y transformación.
Cuándo usarlo: cuando el volumen es bajo y no justificás una capa intermedia, o cuando la transformación es muy simple (renombrar columnas, castear tipos). En la práctica, lo uso poco porque prefiero separar responsabilidades.
6. Compactor
Consolida archivos pequeños en archivos más grandes. Si tenés un pipeline de streaming o micro-batches que genera miles de archivos chicos por hora, el rendimiento de lectura se degrada porque el motor tiene que abrir miles de archivos. El Compactor los junta periódicamente.
Cuándo usarlo: siempre que tengas pipelines de streaming o micro-batches que escriban a Delta, Parquet o Iceberg. Es casi obligatorio en producción.
-- Compactor en Databricks — OPTIMIZE es la implementación nativa
-- Compactar toda la tabla
OPTIMIZE bronze.events;
-- Compactar solo la partición de ayer (mucho más eficiente)
OPTIMIZE bronze.events
WHERE event_date = current_date() - INTERVAL 1 DAY;
-- Si usás Liquid Clustering, OPTIMIZE es incremental
-- y además reordena los datos para queries más rápidas7. Readiness Marker
Señaliza cuándo un dataset está listo para consumo downstream. Parece un detalle menor, pero en pipelines complejos con muchas dependencias es fundamental. Sin un marcador de readiness, tus jobs downstream no saben si los datos ya están completos o si la ingesta sigue corriendo.
Cuándo usarlo: orquestación de pipelines con dependencias complejas. Se puede implementar con archivos _SUCCESS, registros en una tabla de metadatos, o señales en el orquestador (Airflow sensors, Databricks task dependencies).
8. External Trigger
La ingesta no corre en un schedule fijo: se dispara por un evento externo. Un archivo llega a S3, un webhook se activa, un mensaje aparece en una cola. Es el patrón para arquitecturas event-driven.
Cuándo usarlo: cuando la fuente no tiene un schedule predecible o cuando necesitás reaccionar en tiempo real. Ejemplos: archivos que un proveedor sube a un bucket, webhooks de APIs, mensajes en Kafka/EventHub.
Mapeo a la arquitectura Medallion
Uno de los ejercicios más útiles que podés hacer es mapear estos patrones a las capas de una arquitectura Medallion (Bronze / Silver / Gold). No es una relación uno a uno, pero hay afinidades claras:
| Capa | Patrones típicos | Rol |
|---|---|---|
| Bronze | Full Loader, Incremental Loader, CDC, Passthrough Replicator, External Trigger | Ingesta cruda. El objetivo es tener una copia fiel de las fuentes con mínima transformación. |
| Silver | Transformation Replicator, CDC (MERGE), Compactor, Readiness Marker | Limpieza y conformado. Acá aplicás deduplicación, tipado, joins y lógica de negocio liviana. |
| Gold | Compactor, Readiness Marker | Agregaciones y modelos de consumo. Tablas optimizadas para dashboards y ML. |
El Compactor y el Readiness Marker son transversales: los necesitás en cualquier capa donde haya escritura frecuente o dependencias downstream.
Comparación rápida
| Patrón | Complejidad | Latencia | Cuándo usarlo |
|---|---|---|---|
| Full Loader | Baja | Alta | Tablas chicas, sin columna de control |
| Incremental Loader | Media | Media | Tablas con updated_at, batch diario/horario |
| CDC | Alta | Baja | Fuentes transaccionales, captura de deletes |
| Passthrough Replicator | Baja | Variable | Capas de staging, Bronze raw |
| Transformation Replicator | Media | Variable | Transformaciones simples en ingesta |
| Compactor | Baja | N/A (mantenimiento) | Post-streaming, micro-batches |
| Readiness Marker | Baja | N/A (señalización) | Orquestación con dependencias |
| External Trigger | Media-Alta | Baja | Arquitecturas event-driven |
Opinión personal: lo que más uso en producción
De estos 8 patrones, los que uso en el 90% de mis proyectos en Databricks son:
- Incremental Loader para todo lo batch. Es el patrón por defecto. Si la fuente tiene
updated_at, no hay razón para hacer Full Load. - CDC con MERGE para fuentes transaccionales donde necesito capturar updates y deletes. En Databricks, el
MERGE INTOde Delta Lake hace que implementar esto sea casi trivial. - Compactor (vía
OPTIMIZE) después de cualquier pipeline de streaming o micro-batches. Sin esto, el rendimiento de lectura se degrada rápido. - Readiness Marker implementado como task dependencies en Databricks Workflows. Simple pero efectivo.
El Full Loader lo reservo para tablas de dimensiones chicas donde no vale la pena complicarse. El External Trigger lo uso cuando trabajo con archivos que llegan a S3/ADLS en horarios impredecibles (Auto Loader de Databricks es ideal para esto).
Los patrones que menos uso son Passthrough Replicator y Transformation Replicator como patrones separados, porque en la práctica mi capa Bronze ya cumple el rol del Passthrough y Silver el del Transformation Replicator. Pero tener el nombre del patrón me ayuda a explicar qué hace cada capa cuando hablo con el equipo.
Links
- Escuchar el episodio en Spotify
- Data Engineering Design Patterns — Bartosz Konieczny
- Documentación de MERGE en Delta Lake
- Auto Loader en Databricks