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
| Schema | Contenido | Quién lo puebla |
|---|---|---|
raw | Exclusivo MPRO: tablas crudas desde SQL Server, 1:1 con la fuente | dlt-pipelines/ (dlt) |
raw_staging | Staging interno de dlt para escrituras merge | dlt (automático, no tocar a mano) |
raw_sat | Datos derivados de XMLs CFDI del SAT — fuente distinta a MPRO, namespace separado a propósito | script no-dlt (consulta_xmls_gastos) |
monitoring | Trace de las corridas de dlt (duración por paso, filas por tabla, éxito/error) — para dashboards, no para análisis de negocio | dlt-pipelines/log_run_metrics.py |
public / scratch | vací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")
| Tabla | Filas (snapshot) | Fuente | Primary key | Cursor | Host |
|---|---|---|---|---|---|
producto | 27,426 | Producto | (pr_cve_producto) | Fecha_Ult_Modif | .207 |
comprobante_digital | 686,603 | Comprobante_Digital | (cd_tabla, cd_documento) | Fecha_Ult_Modif | .207 |
compra | 149,797 | Compra | (co_folio, co_id) | Fecha_Ult_Modif | .207 |
compra_encabezado | 69,095 | Compra_Encabezado | (co_folio) | Fecha_Ult_Modif | .207 |
orden_compra | 121,985 | Orden_Compra | (oc_folio, oc_id) | Fecha_Ult_Modif | .207 |
existencia | 84,704 | Existencia | (sc_cve_sucursal, al_cve_almacen, pr_cve_producto, tl_cve_talla, cl_cve_color) | Fecha_Ult_Modif | .207 |
movimiento | 4,529,150 | Movimiento | (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 consql_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 resourcepymssqlpuro en vez desql_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
movimientoa.207: el backfill por año usó un filtro manual deMv_Fecha, sindlt.sources.incremental(...)— nunca quedó guardado un cursor deFecha_Ult_Modifen el estado del pipeline. Al correr el incremental por primera vez coninitial_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_Modifsí existe). El fix fue sembrarinitial_valuea mano con elMAX(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 usadlt.sources.incremental(...), el primer corte a incremental necesita sembrarinitial_valuea mano con el máximo cursor ya cargado, o dlt intenta re-traer todo el histórico contra la base viva.
Cron — resumen completo
| Hora | Qué corre | Script |
|---|---|---|
| 5:00 AM | raw_sat.cfdi_recibidos | exploracion/consulta_xmls_gastos/cargar_cfdi_recibidos.py |
| 6:00 AM | comprobante_digital | load_comprobante_digital.py → run_incremental_207() |
| 6:15 AM | reorden + producto + 6 catálogos | load_reorden.py (main()) |
| 6:30 AM | compra, compra_encabezado, orden_compra, existencia | load_compras_inventario.py → run_incremental_207_all() |
| 6:45 AM | movimiento | load_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).
| Tabla | Contenido |
|---|---|
pipeline_runs | 1 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_tables | 1 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)
EOFVéase también
- Pipelines dlt — el patrón genérico detrás de todas las tablas incrementales de esta página.
- Conexiones a TRIVASADB — a qué apunta
.204/.207/.200. - Convención de esquemas y nombres — por qué
raw/raw_satestán separados y cómo se nombrarán las capas siguientes (staging/marts).