Ir al contenido
SparquetSparquet

Pipelines

Sparquet ejecuta un JSON a la vez. Una carga real suele ser varios: truncar, después bronze, después silver, después un commit. Lo único que un JSON solo no puede expresar es el orden en que corren varios de ellos — y eso es un Pipeline.

Un Pipeline es un conjunto ordenado de Jobs del mismo Workflow. No guarda JSON propio: cada etapa guarda una referencia a un Job, y el JSON se compila del canvas de ese Job al momento de ejecutar. Así una etapa nunca se desfasa del archivo que nombra, y un Job borrado deja una etapa que se reporta como rota en vez de una copia vieja que igual corre.

  1. Créalo. Dentro de un Workflow, pulsa New pipeline y ponle el nombre de la carga que ejecuta — Registro nocturno, no secuencia 2. Se abre en /pipelines/:id.

  2. Añade las etapas. Los Jobs del Workflow están listados a la izquierda. Haz clic en uno para añadirlo al final de la fila, o arrástralo al canvas para colocarlo donde lo sueltes. El mismo Job puede aparecer más de una vez — correr el mismo archivo dos veces sobre datos distintos es legítimo — así que nada desaparece de la lista tras usarlo; muestra cuántas veces cada Job ya está en escena.

  3. Dibuja el orden. Conecta el handle derecho de una caja al handle izquierdo de la siguiente. Aquí nada se infiere de los paths, a propósito: dos etapas que no comparten ningún path pueden igual tener que correr en un orden fijo — un truncate antes de una carga — así que el orden es el que dibujas.

  4. Entra en la caja cuando haga falta. Abrir una caja te lleva al canvas del propio Job, con todos los paneles que normalmente tiene: inspector, JSON, issues, IA. Vuelve y la etapa ya refleja la edición, porque la etapa nunca fue una copia.

Cada caja muestra el número de la etapa en el orden de ejecución, lo que el Job lee, lo que escribe, cuántas transformaciones aplica y si tiene bloque de validations — lo suficiente para reconocer una etapa sin abrirla.

Los links que dibujas lo definen; los empates se rompen alfabéticamente, así que la numeración es estable entre renders y entre máquinas.

  • Un ciclo se rechaza mientras lo dibujas: una secuencia con un loop no tiene primera etapa, así que no podría correr.
  • Una etapa sin ningún link igual corre, al final, y se marca como advertencia — para que su posición en la secuencia sea deliberada y no un accidente de cuándo soltaste la caja.

Cómo una etapa le entrega datos a la siguiente

Sección titulada «Cómo una etapa le entrega datos a la siguiente»

Las etapas no se pasan un DataFrame entre sí. Comparten una sesión Spark, y una etapa simplemente lee lo que una anterior escribió:

  • Un path o una tabla. La etapa 1 escribe bronze.pedidos; la etapa 2 lee bronze.pedidos.
  • Una temp view. La etapa 1 escribe una salida view; la etapa 2 lee esa view como su input. Nada toca el almacenamiento — la view vive en la sesión que ambas comparten.

No hay cableado extra para esto, ni hace falta: un Job que ya lee el lugar correcto ya está conectado. Un link en el canvas define cuándo corre una etapa, no qué recibe.

El mismo runner local y el mismo token que un Job suelto. El Pipeline postea todas las etapas en orden a POST /run/flow/stream, que transmite Server-Sent Events, así que el panel se va llenando mientras Spark trabaja en vez de quedarse congelado hasta el final.

Qué obtienes Detalle
Estado por etapa cada caja pasa a running, succeeded, skipped o failed según avanza la secuencia
Resultados por etapa filas leídas, filas escritas, duración y validaciones, para cada etapa
Logs etiquetados por etapa cada línea lleva la etapa que la emitió, incluidas stdout y las líneas de la JVM
Preview hasta 50 filas de la salida de la última etapa — el resultado del Pipeline

La ejecución se detiene en la primera etapa que falla: las siguientes nunca arrancan, y el error nombra la etapa que se rompió. El Pipeline completo toma el mismo lock único de ejecución que un Job, así que una segunda ejecución iniciada mientras una está en curso recibe 409.

Todo o nada, a propósito: una secuencia solo tiene sentido completa, así que una etapa rota detiene la ejecución en vez de correr el resto en silencio con un orden acortado.

Bloqueo Solución
Una etapa apunta a un Job que ya no existe borra la etapa, o recrea el Job al que referenciaba
El Job de una etapa todavía no compila ábrelo y resuelve las issues bloqueantes en su canvas
Una etapa está en un loop quita uno de los links que cierra el loop

Es una conveniencia de desarrollo, como el runner mismo — sin agenda, sin retries, sin alertas, sin ventana de backfill. En producción, entrega los mismos archivos JSON a Airflow, Dagster, Databricks Workflows o cron, en el mismo orden. Lo que un Pipeline te da es poder correr y leer ese orden localmente, antes de que se convierta en la guardia de las 3 de la mañana de alguien.