Ir al contenido
SparquetSparquet

Cargas incrementales y upserts

Una recarga completa es simple y, a partir de cierto tamaño, imposible. Estos son los patrones que la reemplazan.

{
"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.

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.

[
{ "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 failure

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. merge sobre una clave de negocio estable produce la misma tabla se ejecute una o cinco veces. append no 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.

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.

Un pipeline incremental que hace join de varios orígenes construye un plan profundo. Dos reglas lo mantienen barato:

  • checkpoint después de los joins pesados, antes de collect y 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.
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