flowchart LR
subgraph fuente ["Fuente"]
r1["Recurso: alumnos"]
r2["Recurso: asignaturas"]
r3["Recurso: matriculas"]
end
fuente --> p["Pipeline"]
p --> e["Extraer"] --> n["Normalizar"] --> c["Cargar"]
c --> d[("Destino")]
p <-.->|"estado"| d
%% Paleta por bloque del libro
classDef origen fill:#ececf0,stroke:#868d9c,stroke-width:1.5px,color:#2b2f3a
classDef carga fill:#fae3c8,stroke:#bf7f28,stroke-width:1.5px,color:#573809
classDef almacen fill:#d5e6f5,stroke:#3d7cb0,stroke-width:1.5px,color:#12354e
class r1,r2,r3 origen
class p,e,n,c carga
class d almacen
13 Herramientas de ingesta
Si repasamos lo visto hasta ahora, la lista de cosas que hay que resolver para traer una tabla de un sitio a otro es incómodamente larga: paginar, autenticar, reintentar con espera creciente, respetar los límites de uso, inferir tipos, crear la tabla en destino, adaptarla cuando el esquema cambie, guardar la marca de agua, aplicar el merge y dejar constancia de qué se cargó.
Nada de eso es interesante. Todo eso hay que hacerlo bien. Y lo peor es que hay que hacerlo otra vez con la siguiente fuente, y con la siguiente. Esa repetición es la que da sentido a una categoría entera de herramientas.
13.1 El panorama
Conviene distinguir dos cosas que se confunden: el conector (el código que sabe hablar con Salesforce, con Postgres o con la API de Stripe) y la plataforma (lo que programa, monitoriza y gestiona credenciales). Cada herramienta reparte el peso de forma distinta.
| Herramienta | Naturaleza | Nota |
|---|---|---|
| Singer / Meltano | Especificación abierta de taps y targets + plataforma | Fue el estándar abierto de facto; catálogo amplio pero de calidad desigual |
| Airbyte | Plataforma con interfaz gráfica, conectores en contenedores | Muchísimos conectores, requiere operar la plataforma |
| Fivetran | Servicio gestionado, comercial | Se paga y funciona; coste por volumen que sorprende al crecer |
| Sling | Utilidad de línea de comandos | Excelente para base de datos a base de datos, sin ceremonia |
| dlt | Librería de Python | El proceso es código propio; sin plataforma que operar |
Vamos a centrarnos en dlt por una razón concreta: no es un servicio al que haya que conectarse ni un sistema que haya que desplegar, es una librería que se importa. Eso significa que un proceso de ingesta es un fichero de Python que vive en nuestro repositorio, se revisa en un pull request y se ejecuta donde queramos. Encaja con la idea de que el ingeniero de datos toma prestadas las prácticas del desarrollo de software, y además hace que los ejemplos quepan en una página.
13.2 Los tres conceptos
dlt tiene una API pequeña. Con entender tres nombres es suficiente:
- Recurso (
resource): una función de Python que produce registros, normalmente conyield. Se corresponde con una tabla en destino. - Fuente (
source): una agrupación de recursos que comparten configuración, típicamente todas las tablas de un mismo sistema. - Pipeline (
pipeline): quien coge una fuente, la normaliza y la escribe en un destino, llevando la cuenta del estado entre ejecuciones.
El caso mínimo, para hacerse una idea de la superficie de la API:
import dlt
@dlt.resource(name="alumnos")
def alumnos():
yield {"id_alumno": 1, "nombre": "Iraitz", "apellido": "Montalbán"}
yield {"id_alumno": 2, "nombre": "Javier", "apellido": "Garcia"}
pipeline = dlt.pipeline(
pipeline_name="secretaria",
destination="duckdb",
dataset_name="staging",
)
pipeline.run(alumnos())Con eso ya tenemos una tabla staging.alumnos creada, tipada e insertada. No hemos escrito un CREATE TABLE, y ese es precisamente el punto.
El paso intermedio merece una mención. dlt aplana las estructuras anidadas por su cuenta: un campo con una lista de objetos se convierte en una tabla hija enlazada por una clave generada. Es la respuesta práctica a ese JSON con empresas anidadas que veíamos en el capítulo de arquitecturas, y evita tener que decidir el modelado antes de haber visto los datos.
13.3 Extraer de una base de datos
Para nuestra secretaría académica, lo habitual será leer directamente de la base de datos. dlt trae una fuente preparada que refleja el esquema del origen.
import dlt
from dlt.sources.sql_database import sql_database
fuente = sql_database(
credentials="postgresql://lector@servidor:5432/secretaria",
table_names=["alumnos", "asignaturas", "cursa"],
)
pipeline = dlt.pipeline(
pipeline_name="secretaria",
destination="duckdb",
dataset_name="staging",
)
pipeline.run(fuente, write_disposition="replace")Eso es una carga completa de tres tablas. Para pasar a incremental basta con declarar la columna que hace de marca de agua y la clave sobre la que resolver el merge.
fuente = sql_database(
credentials="postgresql://lector@servidor:5432/secretaria",
table_names=["alumnos"],
).with_resources("alumnos")
fuente.alumnos.apply_hints(
incremental=dlt.sources.incremental("actualizado"),
primary_key="id_alumno",
write_disposition="merge",
)
pipeline.run(fuente)Todo el capítulo anterior (la marca de agua, dónde guardarla, cómo recuperarla en la siguiente ejecución, el merge por clave) cabe en esas tres líneas. El estado se persiste en el propio destino, en una tabla _dlt_pipeline_state, con lo que el proceso puede ejecutarse desde una máquina distinta cada vez sin perder el hilo.
En los ejemplos aparecen cadenas de conexión completas por legibilidad. En un proyecto real van en .dlt/secrets.toml (fuera del control de versiones) o en variables de entorno, y dlt las resuelve sola. Una credencial de producción en un repositorio es un incidente de seguridad, no un descuido.
13.4 Extraer de una API
Aquí es donde la herramienta se gana el sueldo. Toda la lista de peajes que enumerábamos en el capítulo de extracción se declara en lugar de programarse.
import dlt
from dlt.sources.rest_api import rest_api_source
fuente = rest_api_source({
"client": {
"base_url": "https://api.universidad.example/v1/",
"auth": {
"type": "bearer",
"token": dlt.secrets["universidad_token"],
},
"paginator": {
"type": "json_link",
"next_url_path": "paging.next",
},
},
"resource_defaults": {
"write_disposition": "merge",
"primary_key": "id",
},
"resources": [
"asignaturas",
{
"name": "alumnos",
"endpoint": {
"path": "alumnos",
"params": {
"actualizado_desde": {
"type": "incremental",
"cursor_path": "actualizado",
"initial_value": "2026-01-01",
},
},
},
},
],
})
pipeline = dlt.pipeline(
pipeline_name="universidad_api",
destination="duckdb",
dataset_name="staging",
)
pipeline.run(fuente)Merece la pena leer ese bloque despacio, porque cada clave sustituye a un problema real. El paginator recorre las páginas hasta el final. El auth mete el token en cada petición. El incremental traduce la marca de agua a un parámetro de consulta y la actualiza al terminar. Los reintentos con espera creciente ante un 429 o un 503 están puestos por defecto. Escrito a mano, esto son doscientas líneas y varias sorpresas.
13.5 Ficheros
Para el CSV en el SFTP o el Parquet en un bucket, el patrón es el mismo, con la particularidad de que conviene apoyarse en la fecha de modificación del fichero para no releer lo ya procesado.
import dlt
from dlt.sources.filesystem import filesystem, read_csv
recurso = filesystem(
bucket_url="s3://datos-universidad/matriculas/",
file_glob="**/*.csv",
) | read_csv()
recurso.apply_hints(
incremental=dlt.sources.incremental("modification_date"),
)
pipeline = dlt.pipeline(
pipeline_name="matriculas_csv",
destination="duckdb",
dataset_name="staging",
)
pipeline.run(recurso.with_name("matriculas"))13.6 Contratos de esquema
Vimos que la política ante un cambio de esquema no debía ser la misma en todas las capas. dlt permite declararla, y con bastante grano fino: qué hacer ante una tabla nueva, ante una columna nueva y ante un tipo que no encaja.
pipeline.run(
fuente,
schema_contract={
"tables": "evolve", # aceptamos tablas nuevas
"columns": "evolve", # aceptamos columnas nuevas
"data_type": "freeze", # pero un cambio de tipo detiene la carga
},
)Los tres valores posibles son evolve (adaptarse), discard_row o discard_value (descartar lo que no encaja) y freeze (detenerse con error). La combinación de arriba es un punto de partida razonable para una capa de aterrizaje: dejamos entrar todo lo nuevo, porque el objetivo es no perder información, pero nos plantamos si una columna que era numérica empieza a traer texto, porque eso no es crecimiento sino que algo se ha roto río arriba.
13.7 Lo que dlt añade a los datos
Al cargar, dlt incorpora columnas técnicas que son exactamente los metadatos de carga que recomendábamos poner a mano:
_dlt_id, identificador único de la fila._dlt_load_id, identificador del lote de carga, que permite aislar o deshacer una ejecución concreta.
Y mantiene además unas cuantas tablas propias en el destino: _dlt_loads con el registro de cada carga y su resultado, _dlt_version con la historia de versiones del esquema inferido, y _dlt_pipeline_state con el estado incremental. Ese conjunto es, en la práctica, la trazabilidad mínima de la capa de ingesta, y viene puesta sin pedirla.
13.8 Quién lo pone en marcha
Hay una pregunta que el capítulo ha esquivado hasta ahora: el proceso está escrito, pero alguien tiene que ejecutarlo. dlt mueve los datos cuando se le llama, y no tiene opinión sobre cuándo hacerlo.
La primera respuesta es siempre cron, y durante un tiempo funciona. Deja de funcionar el día que aparece cualquiera de estas cuatro cosas, y aparecen todas:
- La transformación tiene que esperar a que la ingesta haya terminado bien, no a que hayan pasado veinte minutos.
- Una carga falla por un corte de red y hay que reintentarla sin levantarse de madrugada.
- Alguien pregunta si el proceso de anteayer se ejecutó, y hay que mirar en un log por SSH para saberlo.
- Se corrige un error y hay que reprocesar los últimos treinta días, en orden.
Un orquestador es la pieza que responde a esas cuatro. No aporta capacidad de cálculo ni sabe nada de datos: aporta dependencias, reintentos, calendario y visibilidad.
13.8.1 dlt ya sabe llamar a dbt
Antes de montar nada conviene saber que la dependencia más común, la de “transforma cuando la carga haya ido bien”, no necesita orquestador. dlt trae un ejecutor de proyectos dbt integrado:
import dlt
pipeline = dlt.pipeline(
pipeline_name="secretaria",
destination="duckdb",
dataset_name="staging",
)
carga = pipeline.run(fuente, write_disposition="replace")
transformacion = dlt.dbt.package(pipeline, "/opt/academia")
modelos = transformacion.run_all()
pruebas = transformacion.test()
for nodo in list(modelos) + list(pruebas):
print(nodo.model_name, nodo.status, nodo.time)Lo interesante no es que ahorre un subprocess. Es que package() recibe el pipeline, y de ahí deduce contra qué base de datos tiene que trabajar dbt: exporta las credenciales del destino al entorno (con el prefijo DLT__) y trae perfiles de serie para cada destino, de modo que un proyecto cuyo perfil se llame como el destino no necesita profiles.yml propio. Si el proyecto ya trae el suyo, dbt lo resuelve como siempre y se puede forzar con package_profile_name. En cualquiera de los dos casos desaparece la posibilidad de que la carga escriba en un sitio y la transformación lea de otro, que es un fallo tan tonto como habitual.
Y si pipeline.run() levanta una excepción, la línea siguiente no se ejecuta. La dependencia está expresada por el propio flujo del programa, sin grafo que declarar.
run_all no ejecuta las pruebas
El nombre invita a pensar que hace todo el trabajo, y no. run_all() encadena dbt deps, dbt seed y dbt run: es el equivalente de dbt run, no de dbt build, y las pruebas de los modelos se quedan fuera. Quien sustituya un dbt build por esta llamada se queda sin sus comprobaciones y sin ningún aviso de que ha ocurrido. Por eso arriba se pide test() aparte y se juntan los dos resultados antes de dar el proceso por bueno.
Así que la respuesta honesta a “¿hace falta un orquestador?” para un proceso de un solo origen es no. Un script y una entrada en cron cubren el caso, y montar Airflow para dos pasos es de las decisiones que se lamentan a los seis meses.
13.8.2 Dónde sigue haciendo falta
Lo que un script no colapsa es la convergencia. En cuanto la transformación depende de varios orígenes que van a su ritmo (la secretaría de noche, la plataforma en línea cada hora, un fichero de tarifas cuando administración se acuerda), la pregunta deja de ser “¿terminó la carga?” y pasa a ser “¿terminaron todas?”. Eso ya no cabe en el flujo de un programa:
from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.bash import BashOperator
with DAG(
"academia_diaria",
start_date=datetime(2025, 9, 1),
schedule="0 3 * * *",
catchup=False,
default_args={
"retries": 2,
"retry_delay": timedelta(minutes=5),
},
) as dag:
origenes = [
BashOperator(
task_id=f"ingesta_{nombre}",
bash_command=f"python /opt/pipelines/{nombre}.py",
)
for nombre in ("secretaria", "plataforma", "tarifas")
]
transformacion = BashOperator(
task_id="transformacion",
bash_command="dbt build --project-dir /opt/academia",
)
origenes >> transformacionLas tres cargas se ejecutan en paralelo, cada una reintenta por su cuenta y la transformación arranca cuando todas han terminado bien. Junto a esa convergencia, hay otras tres cosas que el script tampoco da y que rara vez se echan de menos hasta que hacen falta:
- Reproceso. Recalcular los últimos treinta días, en orden y sin solaparse, es una orden en un orquestador y un bucle escrito a mano en cualquier otro sitio.
- Historia. Saber si el proceso de anteayer se ejecutó y cuánto tardó, sin entrar por SSH a leer un log.
- Aviso. Que alguien se entere de que ha fallado sin tener que mirar.
Conviene fijarse en lo que no hay ahí dentro. No hay una sola línea de SQL, ni una transformación, ni una regla de negocio: solo llamadas a procesos que ya existían y saben ejecutarse solos. Esa es la regla que más se incumple y la que más cara sale, porque en cuanto la lógica se muda al orquestador deja de poder ejecutarse en local, deja de poder probarse y queda atada a la herramienta.
El otro detalle es que la validación no necesita una tarea propia. dbt build ejecuta los modelos y sus pruebas en el mismo paso, y se detiene si alguna falla, de modo que la tarea siguiente no llega a arrancar. Separar “transformar” de “validar” en dos tareas suena más ordenado y en realidad abre una ventana para que algo lea datos que ya se sabían malos.
13.8.3 Cuál elegir
| Herramienta | Enfoque | Encaja cuando |
|---|---|---|
| Airflow | Tareas y dependencias, el estándar de facto | Hay equipo de plataforma y se valora encontrar respuestas en internet |
| Dagster | Orientado a los activos de datos, no a las tareas | El proyecto gira alrededor de dbt y se quiere linaje de serie |
| Prefect | Python corriente, decoradores sobre funciones | El equipo es pequeño y ya piensa en Python |
| Kestra | Declarativo en YAML | Se prefiere configuración a código y perfiles mixtos |
La diferencia importante no está en la tabla sino en el enfoque de Dagster: declarar qué activos existen y cómo se derivan unos de otros, en lugar de qué tareas se ejecutan en qué orden. Sobre un proyecto dbt eso encaja casi sin fricción, porque el grafo ya está descrito en los ref().
Un orquestador reintenta, y reintentar es exactamente lo que rompe un proceso que no tolera ejecutarse dos veces. La idempotencia no es un requisito que traiga el orquestador, es uno que el orquestador da por supuesto: la primera vez que una tarea se reintente sola, se sabrá si estaba bien escrita.
Hay un ejercicio completo con este montaje en el apéndice del lago de datos, donde el proceso son dos órdenes seguidas y se explica por qué ahí no se llegó a montar un orquestador: con dos pasos que se lanzan a mano, la maquinaria cuesta más de lo que ordena.
13.9 Lo que una herramienta de ingesta no hace
Para cerrar con expectativas ajustadas. dlt, o cualquiera de sus alternativas, resuelve el movimiento de datos. No decide qué es una entidad de negocio, no construye el modelo dimensional, no calcula métricas y no sabe si el dato es correcto. Todo eso es la transformación, y sigue siendo el trabajo con criterio.
Lo que sí ha cambiado es que la parte mecánica ya no justifica un proyecto de seis meses. Con lo visto en este capítulo tenemos los datos saliendo del origen; queda decidir dónde aterrizan y con qué garantías, que es lo que construiremos a continuación con DuckLake.