Warehouse Postgres trivasa_dw

Réplica en Postgres de partes de TRIVASADB, mantenida por pipelines dlt (ver Pipelines dlt) más un par de cargas ad-hoc no-dlt. Conexión: localhost:5433, db trivasa_dw, user trivasa (credenciales en dlt-pipelines/.dlt/secrets.toml, sección [destination.postgres.credentials]).

Este inventario es un snapshot generado consultando information_schema directamente (2026-07-29) — si cambia algo (nueva tabla, nuevo schema), hay que volver a correr las queries y actualizar, no asumir que sigue vigente indefinidamente.

Schemas

SchemaContenidoQuién lo puebla
rawExclusivo MPRO: tablas crudas desde SQL Server, 1:1 con la fuentedlt-pipelines/ (dlt)
raw_stagingStaging interno de dlt para escrituras mergedlt (automático, no tocar a mano)
raw_satDatos derivados de XMLs CFDI del SAT — fuente distinta a MPRO, namespace separado a propósitoscript no-dlt (consulta_xmls_gastos)
monitoringTrace de las corridas de dlt (duración por paso, filas por tabla, éxito/error) — para dashboards, no para análisis de negociodlt-pipelines/log_run_metrics.py
public / scratchvacíos—

raw — catálogos (write_disposition="replace", recarga completa cada corrida)

familia, sub_familia, categoria, departamento, almacen, sucursal (desde .204) + proveedor (desde .204) + reorden (desde .207, 2,022 filas).

reorden se cambió de merge+incremental a replace porque merge dejaba filas huérfanas cuando se borraban registros en el origen (dlt nunca borra en destino con merge, solo hace upsert) — confirmado: con merge Postgres tenía 2,053 filas vs 2,022 reales en .207 (31 huérfanas); con replace quedó exacto. Es chica (~2k filas, carga completa ~18s junto con producto y los 6 catálogos), así que recargar completo cada vez es más simple y correcto que mantener el upsert.

raw — incrementales (write_disposition="merge")

TablaFilas (snapshot)FuentePrimary keyCursorHost
producto27,426Producto(pr_cve_producto)Fecha_Ult_Modif.207
comprobante_digital686,603Comprobante_Digital(cd_tabla, cd_documento)Fecha_Ult_Modif.207
compra149,797Compra(co_folio, co_id)Fecha_Ult_Modif.207
compra_encabezado69,095Compra_Encabezado(co_folio)Fecha_Ult_Modif.207
orden_compra121,985Orden_Compra(oc_folio, oc_id)Fecha_Ult_Modif.207
existencia84,704Existencia(sc_cve_sucursal, al_cve_almacen, pr_cve_producto, tl_cve_talla, cl_cve_color)Fecha_Ult_Modif.207
movimiento4,529,150Movimiento(mv_folio, mv_id)Fecha_Ult_Modif.207

comprobante_digital: histórico completo 2012-01-10 en adelante, 14 valores de cd_tabla distintos (FACTURA, COMPRA, GASTO_REGISTRO, NOMINA, etc.). Detalle específico del pipeline en Pipelines dlt. Consumido filtrando cd_tabla='COMPRA' para reemplazar un CSV estático de compras en Consulta XMLs vs Gastos/Compras.

Dos narrativas de ingeniería que vale la pena no repetir

  • movimiento — backfill trozado por año, no de un tirón. La carga completa con sql_table/SQLAlchemy (4.4M filas) moría sin traceback: VM de 5.2GB con swap al límite, background killed por falta de memoria. Se resolvió trozando el backfill por año (Mv_Fecha, 2017–2026) con un resource pymssql puro en vez de sql_table. Para una tabla nueva de volumen similar, preferir este patrón desde el inicio en vez de descubrir el límite de memoria en producción.
  • Gotcha del cutover de movimiento a .207: el backfill por año usó un filtro manual de Mv_Fecha, sin dlt.sources.incremental(...) — nunca quedó guardado un cursor de Fecha_Ult_Modif en el estado del pipeline. Al correr el incremental por primera vez con initial_value=1900-01-01, dlt no tenía de dónde retomar y trató de volver a traer las 4.4M filas desde .207 (>20 min, proceso matado). No era problema de índice (IX_Movimiento_Fecha_Ult_Modif sí existe). El fix fue sembrar initial_value a mano con el MAX(Fecha_Ult_Modif) real ya cargado en Postgres para la primera corrida incremental — después de eso, dlt retoma solo del estado persistido. Regla general: si un backfill no usa dlt.sources.incremental(...), el primer corte a incremental necesita sembrar initial_value a mano con el máximo cursor ya cargado, o dlt intenta re-traer todo el histórico contra la base viva.

Cron — resumen completo

HoraQué correScript
5:00 AMraw_sat.cfdi_recibidosexploracion/consulta_xmls_gastos/cargar_cfdi_recibidos.py
6:00 AMcomprobante_digitalload_comprobante_digital.py → run_incremental_207()
6:15 AMreorden + producto + 6 catálogosload_reorden.py (main())
6:30 AMcompra, compra_encabezado, orden_compra, existenciaload_compras_inventario.py → run_incremental_207_all()
6:45 AMmovimientoload_movimiento.py → run_incremental_207()

Todas las tablas de raw que se pueden mantener al día solas están en cron.

raw_staging

Staging transitorio de dlt para resolver el merge de comprobante_digital — no es fuente confiable de datos completos, ignorar para consultas.

raw_sat — fuente SAT, fuera del alcance de dlt

cfdi_recibidos (10,391 filas) — XMLs CFDI parseados desde una carpeta de red montada vía CIFS (//192.168.117.211/SincronizarXml), sin pasar por dlt (carga directa psycopg2/pandas, upsert por uuid, cron diario 5:00 AM). Se separó de un schema raw que en algún punto compartía con MPRO — raw es exclusivo de datos vía dlt desde TRIVASADB; esto viene de otra fuente y necesitaba su propio namespace. Migración fue un ALTER TABLE ... SET SCHEMA, sin re-parsear XMLs.

monitoring — trace de dlt para dashboards (2026-08-06)

Ninguno de los mecanismos nativos de dlt sirve para un dashboard: pipeline.last_trace solo vive en memoria y se pickelea a un archivo local (.dlt/pipelines/<pipeline_name>/trace.pickle) que la siguiente corrida sobreescribe, y raw._dlt_loads (bookkeeping interno de dlt) no tiene ni filas cargadas por tabla ni desglose de duración por paso. Ver Pipelines dlt § trace de corridas para el mecanismo completo (log_run_metrics.run_tracked, ya conectado en los 4 scripts de cron).

TablaContenido
pipeline_runs1 fila por pipeline.run(): pipeline_name, script (qué función lo llamó), load_id, started_at/finished_at, duration_seconds, extract_seconds/normalize_seconds/load_seconds, status (success/failed), error_message
pipeline_run_tables1 fila por tabla tocada en esa corrida (run_id, table_name, rows_count) — viene de trace.last_normalize_info.row_counts

Query base para un dashboard de duración/éxito por pipeline en el tiempo:

select pipeline_name, script, started_at, duration_seconds, status
from monitoring.pipeline_runs
order by started_at desc;

Si el INSERT a monitoring falla, el script solo imprime el error — nunca tumba la carga real (el try/except de log_run_metrics.py envuelve únicamente el logging, no el pipeline.run()).

Cómo regenerar este inventario

cd dlt-pipelines && source venv/bin/activate
python3 - <<'EOF'
import psycopg2
conn = psycopg2.connect(host='localhost', port=5433, dbname='trivasa_dw', user='trivasa', password='<ver .dlt/secrets.toml>')
cur = conn.cursor()
cur.execute("""
    SELECT table_schema, table_name
    FROM information_schema.tables
    WHERE table_schema NOT IN ('pg_catalog','information_schema')
    ORDER BY 1,2
""")
for row in cur.fetchall():
    print(row)
EOF

Véase también