Pipelines dlt — patrón y runbook

Cómo se replica una tabla de SQL Server (TRIVASADB) a Postgres (trivasa_dw) con dlt, y cómo armar un pipeline nuevo siguiendo el mismo patrón. Vive en dlt-pipelines/ del repo trivasa-bi-dev. El pipeline de referencia es load_comprobante_digital.py; load_reorden.py sigue una variante del mismo patrón.

⚠️ Manejo de credenciales: las credenciales del destino Postgres viven en .dlt/secrets.toml ([destination.postgres.credentials]), fuera de git. Las credenciales del origen SQL Server, en cambio, están hardcodeadas directamente en los scripts de pipeline (no en secrets.toml) — un patrón que conviene corregir antes de que el repo tenga más colaboradores. Esta página describe el mecanismo sin reproducir los valores reales; verlos en el script correspondiente si se tiene acceso al repo.

Entorno

cd dlt-pipelines
python3 -m venv venv
source venv/bin/activate
pip install "dlt[postgres]" pymssql

Sin problemas de compilación (pymssql, psycopg2-binary, orjson bajan como wheels precompilados). Nota curiosa sin impacto: un venv nuevo bajo Python 3.14 incluye un symlink venv/bin/𝜋thon — easter egg real de CPython 3.14, no una alteración del sistema.

Servidores SQL Server: cuál usar para qué

Ver Conexiones a TRIVASADB para la tabla completa. Resumen para este patrón: .200/TRIVASADB3 para la carga inicial (backfill masivo, no pega a producción), .207/TRIVASADB para la sincronización incremental continua (base viva). .200/TRIVASADB (sin el 3) nunca se usa — está confirmada desactualizada.

El patrón: un resource parametrizado por conexión

def comprobante_digital(server, user, password, database):
    @dlt.resource(name="comprobante_digital", write_disposition="merge", primary_key=["Cd_Tabla", "Cd_Documento"])
    def _resource(fecha=dlt.sources.incremental("Fecha_Ult_Modif", initial_value=datetime.datetime(1900, 1, 1))):
        conn = pymssql.connect(server=server, user=user, password=password, database=database)
        cursor = conn.cursor(as_dict=True)
        cursor.execute("""
            SELECT Cd_Tabla, Cd_Documento, Cd_Timbre_UUID, Fecha_Alta, Fecha_Ult_Modif
            FROM Comprobante_Digital
            WHERE Fecha_Ult_Modif > %s
        """, (fecha.last_value,))
        for row in cursor:
            yield row
        conn.close()
    return _resource

La clave: server/user/password/database son parámetros de la función, no del decorador @dlt.resource. Esto permite llamar al mismo resource contra dos conexiones distintas (.200 para el backfill, .207 para incremental) sin duplicar código — usando el mismo pipeline_name, para que el estado incremental de dlt (el last_value de Fecha_Ult_Modif) se comparta entre ambas corridas.

Gotcha confirmado: initial_value debe ser datetime.datetime(1900, 1, 1), no el string "1900-01-01". Si la columna cursor es datetime en SQL Server, dlt intenta comparar str > datetime y falla con IncrementalCursorInvalidCoercion. load_reorden.py todavía tiene el bug con el string sin corregir — replicarlo ahí si se vuelve a tocar ese pipeline.

Paso 1 — Carga inicial desde .200/TRIVASADB3

if __name__ == "__main__":
    pipeline = dlt.pipeline(pipeline_name="comprobante_digital_compra", destination="postgres", dataset_name="raw")
    resource = comprobante_digital("192.168.117.200", "<usuario .200>", "<password .200>", "TRIVASADB3")
    print(pipeline.run(resource()))

Es una carga masiva de todo el histórico (initial_value = 1900); hacerle un full-scan sin filtro reciente a .207 (producción en vivo) es más riesgo que usar la copia .200/TRIVASADB3 para absorber ese costo. Al terminar, dlt guarda el last_value alcanzado de Fecha_Ult_Modif en el estado local del pipeline — eso es lo que hace posible el corte a incremental del paso 3.

Paso 2 — Elegir el primary_key de sincronización

write_disposition="merge" necesita un primary_key que identifique de forma única cada fila, o se duplican/pierden datos en el upsert. Validar contra ambas bases (backfill y en vivo pueden tener distinto nivel de “suciedad”) con:

SELECT Cd_Tabla, Cd_Documento, COUNT(*)
FROM Comprobante_Digital
GROUP BY Cd_Tabla, Cd_Documento
HAVING COUNT(*) > 1

Si no regresa filas en ninguna de las dos bases, el PK es válido. No asumir que la PK “obvia” de negocio alcanza — para Reorden se necesitó una PK compuesta de 6 columnas por duplicados reales de talla/color (ver Modelo de datos TRIVASADB).

Paso 3 — Cambiar a la BD en vivo .207

def run_incremental_207():
    pipeline = dlt.pipeline(pipeline_name="comprobante_digital_compra", destination="postgres", dataset_name="raw")
    resource = comprobante_digital("192.168.117.207", "<usuario .207>", "<password .207>", "TRIVASADB")
    print(pipeline.run(resource()))

Como reusa el mismo pipeline_name, dlt retoma el last_value de Fecha_Ult_Modif dejado por el backfill de .200 y solo trae de .207 las filas más nuevas que ese punto de corte — no repite el histórico completo contra la base viva.

Riesgo a validar antes de cortar a .207 en un pipeline nuevo: el cutover solo es seguro si .200/TRIVASADB3 no tiene un Fecha_Ult_Modif “atrasado” respecto a .207 para ninguna fila relevante. Si .200/TRIVASADB3 está desfasada (no confirmado al 100% que siempre esté al día), el last_value heredado podría dejar huecos — para un pipeline crítico, comparar manualmente el MAX(Fecha_Ult_Modif) de ambas bases antes de depender del cutover automático.

Cómo replicar el patrón para una tabla nueva

  1. Copiar la función server/user/password/database → @dlt.resource con incremental(..., initial_value=datetime.datetime(1900,1,1)).
  2. Definir el SELECT ... WHERE <col_fecha> > %s con la columna de auditoría real de la tabla (no siempre se llama Fecha_Ult_Modif — confirmar con Modelo de datos TRIVASADB o el dump docs/trivasadb_schema.txt del repo).
  3. Validar el primary_key con la query de duplicados del paso 2, en .200/TRIVASADB3 y en .207.
  4. Backfill inicial contra .200/TRIVASADB3 (mismo pipeline_name para todo el ciclo de vida de esa tabla).
  5. Confirmar que .200/TRIVASADB3 no está desfasada respecto a .207 para el rango de fechas relevante.
  6. Agregar la función de corte a .207 reusando el mismo pipeline_name, y programarla en cron para las corridas incrementales continuas.

Si la tabla es grande (millones de filas) y el backfill con sql_table/SQLAlchemy corre el riesgo de quedarse sin memoria, ver la narrativa de movimiento (trozado por año) en Warehouse Postgres antes de intentar una sola corrida completa.

Trace de corridas: log_run_metrics.py

dlt no persiste el trace de sus corridas por su cuenta: pipeline.last_trace solo vive en memoria y, al terminar el run, se pickelea a .dlt/pipelines/<pipeline_name>/trace.pickle — un archivo local que la siguiente corrida sobreescribe. raw._dlt_loads (bookkeeping interno de dlt en Postgres) tampoco alcanza para un dashboard: no tiene filas cargadas por tabla ni desglose de duración por paso (extract/normalize/load).

dlt-pipelines/log_run_metrics.py cierra ese hueco. Cada script llama log_run_metrics.run_tracked(pipeline, data, script) en vez de pipeline.run(data) directo — corre el pipeline.run(...) normal y, tanto si tiene éxito como si falla, inserta el trace (pipeline.last_trace) en monitoring.pipeline_runs / monitoring.pipeline_run_tables (ver Warehouse Postgres, sección monitoring). Si el INSERT a monitoring falla, solo se imprime el error — el try/except envuelve únicamente el logging, nunca el pipeline.run() real, para que un problema de trace no pueda tumbar una carga.

Ya conectado en los 4 scripts de cron (load_reorden.py, load_movimiento.py, load_compras_inventario.py, load_comprobante_digital.py). Para un pipeline nuevo: importar log_run_metrics y reemplazar pipeline.run(data) por log_run_metrics.run_tracked(pipeline, data, "script.py:funcion") — no cambia el comportamiento de la carga, solo agrega el registro del trace.

Pipeline en producción: raw.comprobante_digital

Fuente: Comprobante_Digital, tabla genérica de comprobantes que cuelga de cualquier tipo de documento vía Cd_Tabla/Cd_Documento. Script: load_comprobante_digital.py.

  • Backfill inicial (.200/TRIVASADB3): 684,466 filas totales, 61,563 de cd_tabla='COMPRA' en ese momento.
  • Cron diario, 6:00 AM, incremental contra .207.
  • Destino: raw.comprobante_digital, columnas normalizadas a snake_case por dlt, PK (cd_tabla, cd_documento).
  • Consumo confirmado: filtrando cd_tabla='COMPRA', reemplazó un CSV estático (data/uuid_compras.csv, 16,346 filas sin refresco) en Consulta XMLs vs Gastos/Compras — la fuente Postgres llega a 61,710 filas y se actualiza a diario. Impacto medido en “pendientes” del mes en curso: bajó de 569 a 562 (julio 2026); el mes anterior no cambió, ya estaba completo en ambas fuentes.

Pendiente

Replicar el mismo patrón para Pago_Cxp_Comprobante (lado gastos, hoy solo disponible como CSV estático uuid_gastos.csv) — mismo consumidor (Consulta XMLs vs Gastos/Compras). Cuando se migre, la lógica que hoy lee ese CSV necesitará el mismo ajuste que se le hizo al de compras.

Véase también