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 ensecrets.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]" pymssqlSin 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 _resourceLa 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(*) > 1Si 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
- Copiar la función
server/user/password/database→@dlt.resourceconincremental(..., initial_value=datetime.datetime(1900,1,1)). - Definir el
SELECT ... WHERE <col_fecha> > %scon la columna de auditoría real de la tabla (no siempre se llamaFecha_Ult_Modif— confirmar con Modelo de datos TRIVASADB o el dumpdocs/trivasadb_schema.txtdel repo). - Validar el
primary_keycon la query de duplicados del paso 2, en.200/TRIVASADB3y en.207. - Backfill inicial contra
.200/TRIVASADB3(mismopipeline_namepara todo el ciclo de vida de esa tabla). - Confirmar que
.200/TRIVASADB3no está desfasada respecto a.207para el rango de fechas relevante. - Agregar la función de corte a
.207reusando el mismopipeline_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 decd_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
- Conexiones a TRIVASADB
- Warehouse Postgres — inventario de todas las tablas cargadas con este patrón, más el cron completo.
- Convención de esquemas