Cargas incrementales y upserts
Una recarga completa es simple y, a partir de cierto tamaño, imposible. Estos son los patrones que la reemplazan.
Merge en lugar de overwrite
Sección titulada «Merge en lugar de overwrite»{ "output": { "format": "delta", "path": "analytics.orders", "mode": "merge", "options": { "merge_keys": ["order_id"], "merge_condition": "T.deleted = false" } }}T es la tabla destino, S el DataFrame entrante. merge_condition agrega SQL a la cláusula de coincidencia — útil para ignorar filas con borrado lógico (soft-delete) o para restringir el merge a una partición.
Lee solo la ventana nueva
Sección titulada «Lee solo la ventana nueva»Empuja la ventana hacia el origen para que la lectura sea pequeña:
// a partitioned table or a database query{ "type": "filter", "condition": "ordered_at >= '{since}'" }fw.run("orders.json", params={"since": last_success_timestamp})Mantener since fuera del archivo — en el orquestador, en una tabla de control — es lo que hace que el pipeline sea reutilizable tanto para el delta nocturno como para un backfill.
Detente cuando no hay nada que hacer
Sección titulada «Detente cuando no hay nada que hacer»[ { "type": "filter", "condition": "updated_at > '{since}'" }, { "type": "stop_if_empty", "message": "No changes since {since}" }, { "type": "checkpoint" }]Las transformaciones restantes y toda escritura se omiten, y la ejecución vuelve con skipped: true y success: true. Colocado justo después del filtro de ventana — antes de los joins y la construcción de payloads — convierte una noche sin actividad en unos pocos segundos en lugar de un pipeline completo sobre cero filas.
result = fw.run("orders.json", params={"since": since})if result.skipped: log.info("nothing to load") # not a failureHaz seguras las reejecuciones
Sección titulada «Haz seguras las reejecuciones»Un job incremental será reejecutado — tras un fallo, tras un archivo que llegó tarde, por alguien depurando. Dos propiedades hacen que eso sea inofensivo:
- Escrituras idempotentes.
mergesobre una clave de negocio estable produce la misma tabla se ejecute una o cinco veces.appendno lo hace. - Una ventana determinista. Derívala de un parámetro o de una tabla de control, nunca de
now()dentro del pipeline, para que una reejecución cubra el mismo rango.
{ "type": "drop_duplicates", "columns": ["order_id"] }Deduplicar sobre la clave de merge justo antes de la escritura cierra el último hueco: un origen que entrega el mismo registro dos veces dentro de una misma ventana.
Staging, y luego commit
Sección titulada «Staging, y luego commit»Cuando varios pipelines contribuyen a un mismo destino, deja que cada uno escriba una vista de staging y dale al final el trabajo de validar y publicar:
// pipelines 1..N{ "output": { "format": "view", "path": "orders_staging", "mode": "overwrite" } }// the commit pipeline{ "input": { "format": "view", "path": "orders_staging" }, "validations": { "on_failure": "fail", "rules": [ { "type": "unique", "columns": ["order_id"] }, { "type": "row_count", "min": 1 } ] }, "outputs": [ { "format": "delta", "path": "analytics.orders", "mode": "merge", "options": { "merge_keys": ["order_id"] } } ]}fw = Sparquet()for conf in ["conf_a.json", "conf_b.json", "conf_c.json"]: result = fw.run(conf, params={"since": since}) if not result.success: raise RuntimeError(f"{conf}: {result.error}")
commit = fw.run("conf_commit.json")fw.stop()Nada llega al destino hasta que todo el conjunto tuvo éxito y las reglas pasaron. Las vistas temporales viven en la sesión de Spark, así que todos los pipelines deben ejecutarse en el mismo proceso — que es exactamente lo que te da reutilizar un único Sparquet.
Vigila el linaje
Sección titulada «Vigila el linaje»Un pipeline incremental que hace join de varios orígenes construye un plan profundo. Dos reglas lo mantienen barato:
checkpointdespués de los joins pesados, antes decollecty antes de abrirse en abanico hacia los destinos.- Empuja la lista de claves hacia las lecturas grandes con pushdown en tiempo de ejecución para que el origen omita archivos en lugar de escanearlos.
Elegir un patrón
Sección titulada «Elegir un patrón»| Situación | Patrón |
|---|---|
| Tabla de dimensión pequeña | overwrite completo — lo simple gana |
| Tabla de hechos con una clave estable | merge sobre la clave |
| Log de eventos solo de append | append más una lectura que deduplica |
| Varios jobs alimentando una tabla | vista de staging, y luego un pipeline de commit |
| Backfill | el mismo archivo, un {since} distinto |