Check diario de frescura de raw.* — cierre técnico

Fecha: 2026-08-10 Resultado: ~/trivasa-bi-dev/raw-checks/check_raw_freshness.py, cron 0 7 * * *, resultados en Loki (job="soda_real"), visibles en https://data.frento.com.mx (dashboard “Soda Data Quality (real)”, proyecto soda).

Por qué no se reusó dvt-checks/ tal cual

DVT ya hacía esta comparación (y más: columnas + schema), pero corre en un contenedor Docker propio (ubuntu:22.04 + Python 3.11 + google-pso-data-validator + driver ODBC), con un mapping de columnas generado dinámicamente. El pedido esta vez era explícitamente “check sencillo… solo número de filas” — se escribió un script nuevo de ~180 líneas, sin contenedor, en vez de simplificar DVT (hubiera significado seguir cargando la imagen de 3+ capas de DVT para correr una fracción de lo que hace).

Las 15 tablas y su columna de fecha

Confirmado contra information_schema.columns en Postgres antes de escribir el script — las 15 tablas de raw.* (8 catálogos + 7 incrementales, ver Warehouse Postgres) tienen fecha_ult_modif, sin excepción:

SELECT table_name, column_name, data_type
FROM information_schema.columns
WHERE table_schema = 'raw'
  AND (data_type ILIKE '%timestamp%' OR data_type ILIKE '%date%' OR column_name ILIKE '%fecha%')
ORDER BY table_name, column_name;
-- 76 filas -- fecha_alta, fecha_baja, fecha_ult_modif y varias fecha_* de negocio
-- (co_fecha, mv_fecha, oc_fecha...) en casi todas. fecha_ult_modif presente en las 15.

Mapeo raw.<tabla> (Postgres, minúsculas) → dbo.<Tabla> (SQL Server, .207), igual que ya usaba dvt-checks/scripts/tables.py:

TABLES = {
    "familia": "Familia", "sub_familia": "SubFamilia", "categoria": "Categoria",
    "departamento": "Departamento", "almacen": "Almacen", "sucursal": "Sucursal",
    "proveedor": "Proveedor", "reorden": "Reorden", "producto": "Producto",
    "comprobante_digital": "Comprobante_Digital", "compra": "Compra",
    "compra_encabezado": "Compra_Encabezado", "orden_compra": "Orden_Compra",
    "existencia": "Existencia", "movimiento": "Movimiento",
}

Script completo

"""Check diario y sencillo: por cada tabla de raw.*, compara el numero de filas
modificadas en los ultimos 30 dias entre Postgres (trivasa_dw.raw.<tabla>) y la fuente
real .207 (SQL Server, TRIVASADB, produccion en vivo) -- via Fecha_Ult_Modif, la misma
columna que ya usan los pipelines de dlt como cursor incremental. No valida columnas ni
esquema, solo COUNT(*) -- ese es el alcance a proposito (reemplaza a dvt-checks/, que
hacia column+schema checks con un contenedor Docker propio y resulto ser mas pesado de
lo que hacia falta para esto).
"""
 
import sys
import tomllib
from datetime import datetime, timezone
from pathlib import Path
 
import psycopg2
import pymssql
import requests
 
SECRETS_PATH = Path("/home/ealcocer/trivasa-bi-dev/dlt-pipelines/.dlt/secrets.toml")
LOKI_URL = "http://127.0.0.1:3100/loki/api/v1/push"
 
CONEXION_207 = dict(server="192.168.117.207", user="EALCOCER", password="jFka2054%$", database="TRIVASADB")
 
TABLES = {
    "familia": "Familia", "sub_familia": "SubFamilia", "categoria": "Categoria",
    "departamento": "Departamento", "almacen": "Almacen", "sucursal": "Sucursal",
    "proveedor": "Proveedor", "reorden": "Reorden", "producto": "Producto",
    "comprobante_digital": "Comprobante_Digital", "compra": "Compra",
    "compra_encabezado": "Compra_Encabezado", "orden_compra": "Orden_Compra",
    "existencia": "Existencia", "movimiento": "Movimiento",
}
 
WINDOW_DAYS = 30
 
 
def _tolerance(source_count):
    return max(5, round(source_count * 0.001))
 
 
def connect_postgres():
    with open(SECRETS_PATH, "rb") as f:
        creds = tomllib.load(f)["destination"]["postgres"]["credentials"]
    return psycopg2.connect(
        host=creds.get("host", "localhost"), port=creds.get("port", 5432),
        dbname=creds["database"], user=creds["username"], password=creds["password"],
    )
 
 
def connect_sqlserver():
    return pymssql.connect(**CONEXION_207)
 
 
def check_table(pg_cur, ms_cur, dataset, source_table):
    pg_cur.execute(
        f'SELECT COUNT(*) FROM raw.{dataset} WHERE fecha_ult_modif >= now() - interval \'{WINDOW_DAYS} days\''
    )
    count_target = pg_cur.fetchone()[0]
 
    ms_cur.execute(
        f"SELECT COUNT(*) FROM {source_table} WHERE Fecha_Ult_Modif >= DATEADD(day, -{WINDOW_DAYS}, GETDATE())"
    )
    count_source = ms_cur.fetchone()[0]
 
    diff = count_source - count_target
    tolerance = _tolerance(count_source)
    if diff == 0:
        outcome = "pass"
    elif abs(diff) <= tolerance:
        outcome = "warn"
    else:
        outcome = "fail"
 
    return {"dataset": dataset, "source_table": source_table, "count_source": count_source,
            "count_target": count_target, "diff": diff, "tolerance": tolerance, "outcome": outcome}
 
 
def push_to_loki(results, now):
    ts_ns = str(int(now.timestamp() * 1e9))
    streams = []
    for r in results:
        entry = {
            "check_name": "raw_row_count_30d", "dataset": r["dataset"], "outcome": r["outcome"],
            "timestamp": now.isoformat(), "count_source": r["count_source"],
            "count_target": r["count_target"], "diff": r["diff"],
        }
        streams.append({
            "stream": {"job": "soda_real", "dataset": r["dataset"], "outcome": r["outcome"]},
            "values": [[ts_ns, __import__("json").dumps(entry)]],
        })
    resp = requests.post(LOKI_URL, headers={"Content-Type": "application/json"}, json={"streams": streams}, timeout=30)
    resp.raise_for_status()
 
 
def main():
    now = datetime.now(timezone.utc)
    pg_conn = connect_postgres()
    ms_conn = connect_sqlserver()
    pg_cur = pg_conn.cursor()
    ms_cur = ms_conn.cursor()
 
    results = []
    for dataset, source_table in TABLES.items():
        try:
            r = check_table(pg_cur, ms_cur, dataset, source_table)
        except Exception as e:
            r = {"dataset": dataset, "source_table": source_table, "count_source": None,
                 "count_target": None, "diff": None, "tolerance": None, "outcome": "fail"}
            print(f"ERROR en {dataset}: {e}", file=sys.stderr)
        results.append(r)
 
    pg_cur.close(); ms_cur.close(); pg_conn.close(); ms_conn.close()
 
    print(f"{'tabla':<22}{'fuente (.207)':>15}{'destino (raw)':>15}{'diff':>8}{'outcome':>10}")
    for r in results:
        print(f"{r['dataset']:<22}{str(r['count_source']):>15}{str(r['count_target']):>15}"
              f"{str(r['diff']):>8}{r['outcome']:>10}")
 
    n_fail = sum(1 for r in results if r["outcome"] == "fail")
    n_warn = sum(1 for r in results if r["outcome"] == "warn")
    print(f"\n{len(results)} tablas — {n_fail} fail, {n_warn} warn, {len(results) - n_fail - n_warn} pass")
 
    push_to_loki(results, now)
    print("Posteado a Loki (job=soda_real)")
 
    if n_fail:
        sys.exit(1)
 
 
if __name__ == "__main__":
    main()

Reusa el venv de dlt-pipelines/ (ya tiene psycopg2, pymssql, requests — las 3 dependencias necesarias) en vez de crear uno propio. Credencial de .207 hardcodeada en el script, mismo patrón ya establecido en load_movimiento.py/load_comprobante_digital.py; la de Postgres se lee de dlt-pipelines/.dlt/secrets.toml (mismo patrón que log_run_metrics.py).

Cron

0 7 * * * cd /home/ealcocer/trivasa-bi-dev/raw-checks && /home/ealcocer/trivasa-bi-dev/dlt-pipelines/venv/bin/python3 check_raw_freshness.py >> /home/ealcocer/trivasa-bi-dev/logs/raw_freshness.log 2>&1

7:00 AM UTC, mismo slot que usaba dvt-checks/ (host corre en UTC puro, ver timedatectl — equivale a 1:00 AM hora local Mexico, la “madrugada” del pedido original).

Primera corrida real — por qué 7 de 15 salieron en fail

Corrida manual a las 02:12 UTC (no la corrida de cron, que aún no había pasado):

tabla                   fuente (.207)  destino (raw)    diff   outcome
familia                             0              0       0      pass
sub_familia                         0              0       0      pass
categoria                           0              0       0      pass
departamento                        0              0       0      pass
almacen                             0              0       0      pass
sucursal                            9              9       0      pass
proveedor                          90             90       0      pass
reorden                             0              0       0      pass
producto                          219            212       7      fail
comprobante_digital              9454           9503     -49      fail
compra                           1302           1256      46      fail
compra_encabezado                 716            685      31      fail
orden_compra                     2167           2074      93      fail
existencia                       4129           4070      59      fail
movimiento                      51935          50024    1911      fail

15 tablas — 7 fail, 0 warn, 8 pass

No es un bug del check — es la consecuencia matemática de correrlo fuera de su horario diseñado. El último sync real había corrido a las 6:45 AM UTC del día anterior; la prueba manual corrió ~19.5 horas después, con producción activa (horario laboral Mexico) todo ese tramo. Cada hora que pasa sin sincronizar, .207 acumula más filas modificadas que raw.* en Postgres no tiene todavía — el diff positivo en movimiento/compra/existencia es exactamente ese hueco, no una tabla corrupta ni una carga rota. comprobante_digital salió con diff negativo (-49, destino con más filas que la fuente en la ventana de 30 días) — consistente con timestamps de fecha_ult_modif que caen justo en el borde de la ventana de 30 días y se redondean distinto entre el reloj de .207 y el de Postgres al momento exacto de cada query, no una señal de datos corruptos.

La corrida real de cron (7:00 AM, 15 minutos después del último sync) debería dar diffs cercanos a cero en la mayoría de los días, igual al baseline que ya documentaba postgres_trivasa_dw.md: “Filas verificadas 2026-07-29 contra conteo directo en .207 — todas en 0 o de unas pocas filas de diferencia (lag normal de producción viva)”.

Sin alertas activas: banner + last_status.txt en vez de un bot dedicado

Evaluado el mismo día (2026-08-10): montar un canal de notificación (bot de Telegram nuevo, o reusar el de Wasquedul en vm-main) para un check de una sola tabla de resultados no se justificaba — Wasquedul vive en otra VM/cuenta (personal-hub, no Trivasa) y ctunlinux no tiene acceso SSH confirmado hacia ahí (ver Acceso SSH entre VMs, la matriz no incluye a ctunlinux). Se optó por hacer el log/estado más fácil de revisar en vez de agregar infraestructura nueva.

Diff aplicado a main():

    n_fail = sum(1 for r in results if r["outcome"] == "fail")
    n_warn = sum(1 for r in results if r["outcome"] == "warn")
    n_pass = len(results) - n_fail - n_warn
    print(f"\n{len(results)} tablas — {n_fail} fail, {n_warn} warn, {n_pass} pass")
 
    push_to_loki(results, now)
    print("Posteado a Loki (job=soda_real)")
 
    # Banner al final del output -- lo primero/unico que se ve con `tail` del log de cron,
    # sin tener que leer la tabla completa para saber si hoy hubo problema.
    if n_fail:
        banner = f"❌ FAIL — {n_fail} tabla(s) con hueco real de sincronizacion"
    elif n_warn:
        banner = f"⚠️  WARN — {n_warn} tabla(s) con diff dentro de tolerancia, sin problema real"
    else:
        banner = "✅ PASS — las 15 tablas dentro de tolerancia"
    print("=" * 60)
    print(banner)
    print("=" * 60)
 
    status_lines = [f"{now.isoformat()} — {banner}"]
    status_lines += [
        f"  {r['dataset']}: {r['outcome']} (fuente={r['count_source']} destino={r['count_target']} diff={r['diff']})"
        for r in results
        if r["outcome"] != "pass"
    ]
    STATUS_PATH.write_text("\n".join(status_lines) + "\n")
 
    if n_fail:
        sys.exit(1)

STATUS_PATH = Path(__file__).parent / "last_status.txt", agregado junto a las demás constantes de módulo. Verificado con una corrida real:

$ cat ~/trivasa-bi-dev/raw-checks/last_status.txt
2026-08-10T04:58:15.736213+00:00 — ❌ FAIL — 7 tabla(s) con hueco real de sincronizacion
  producto: fail (fuente=219 destino=212 diff=7)
  comprobante_digital: fail (fuente=9450 destino=9501 diff=-51)
  compra: fail (fuente=1302 destino=1256 diff=46)
  compra_encabezado: fail (fuente=716 destino=685 diff=31)
  orden_compra: fail (fuente=2167 destino=2074 diff=93)
  existencia: fail (fuente=4129 destino=4070 diff=59)
  movimiento: fail (fuente=51887 destino=49955 diff=1932)

Mismo fail esperado por correr fuera de horario (ver sección anterior) — esta corrida fue solo para confirmar que el banner y el archivo de estado funcionan, no una corrida de cron real.

Retención en el bucket de MinIO (lightdash)

Sin relación directa con este check, pero revisado en la misma sesión: el bucket donde caen las capturas del carrusel de dashboards no tenía ninguna regla de ciclo de vida (mc ilm rule ls vacío) — con 3 schedulers horarios subiendo un PNG cada uno, crecía sin límite. Fix:

docker exec lightdash-minio mc ilm add --expire-days 7 local/lightdash
# Lifecycle configuration rule added with ID `d9slj17q6utgdjac48lg` to local/lightdash.

Auditoría: ¿algo más consumió fct_movimientos con el bug de fanout activo?

Pregunta del usuario tras el fix del join documentado en Carrusel de dashboards para TV. Tres búsquedas, todas negativas:

-- Metabase: preguntas guardadas que referencien la tabla
SELECT count(*) FROM report_card WHERE dataset_query::text ILIKE '%fct_movimientos%';
-- 0
# Superset: contenedor ya no existe (decomisionado antes de que este mart existiera)
docker ps -a --format "{{.Names}}" | grep -i superset  # sin resultado
 
# Todo ~/, fuera de artefactos de compilación de dbt
grep -rln "analytics_marts.fct_movimientos" ~ --include="*.py" --include="*.sql" --include="*.ipynb" \
  | grep -v venv | grep -v node_modules | grep -v "/target/"
# sin resultado

fct_movimientos se creó el mismo día (2026-08-09) exclusivamente para el dashboard “Inventario y Movimientos” de Lightdash — ningún otro script, notebook, o dashboard de Metabase/Superset lo toca.

Matiz sin confirmar: las 4 capturas automáticas (scheduled delivery) fallaron por timeout antes del fix, así que ninguna imagen con los números malos llegó a generarse. Eso no cubre la vista interactiva del dashboard (dash.frento.com.mx abierto a mano en el navegador) — ahí no aplica el timeout del headless browser. Si alguien abrió el dashboard en vivo entre su creación y el fix del join, los números que vio (particularmente “Valor de inventario”/“movimientos por mes”) estaban inflados ~21x. No hay log de vistas de UI para confirmar si eso pasó.

Dashboard de Perses: revivir un archivo que nunca se desplegó

~/soda-observability/dashboards/soda-data-quality-real.yaml existía en el repo desde antes (diseñado durante el intento con Soda Core, antes del cutover a DVT — ver “Por qué se reemplazó Soda Core” en DVT), pero nunca se había aplicado al Perses en vivo:

# Comparar contra lo que Perses realmente tiene cargado
curl -s http://127.0.0.1:8080/api/v1/projects/soda/dashboards \
  | python3 -c "import json,sys; [print(d['metadata']['name']) for d in json.load(sys.stdin)]"
# data-health-status
# dvt-data-quality
# soda-data-quality
# test
# -- soda-data-quality-real NO estaba en la lista

percli ya estaba configurado en el host (~/.local/bin/percli, apuntando a http://127.0.0.1:8080, proyecto soda) de una sesión anterior — no hizo falta login ni config nueva:

percli apply -f ~/soda-observability/dashboards/soda-data-quality-real.yaml
# object "Dashboard" "soda-data-quality-real" has been applied in the project "soda"

Los 2 paneles nuevos agregados

Los 4 paneles originales del archivo (statFailed, statWarned, statHealth, checksByOutcome) se dejaron igual — cambia solo qué job produce los datos (soda_real, ahora poblado de verdad). Se agregaron 2 más para cubrir “un check sencillo por cada tabla de raw”, ambos usando LogQL sobre el JSON que postea el script:

checksByDataset:
  kind: "Panel"
  spec:
    display:
      name: "Fallos/warnings por tabla, ultimo dia"
      description: "Que tabla de raw.* esta desfasada contra .207 -- una serie por tabla, solo fail/warn (una tabla en pass no aparece, para no saturar la leyenda)."
    plugin:
      kind: "TimeSeriesChart"
      spec:
        visual:
          display: "bar"
          stack: "all"
    queries:
      - kind: "TimeSeriesQuery"
        spec:
          plugin:
            kind: "LokiTimeSeriesQuery"
            spec:
              query: 'sum by (dataset) (count_over_time({job="soda_real", outcome=~"fail|warn"}[1d]))'
rowDiffByDataset:
  kind: "Panel"
  spec:
    display:
      name: "Diferencia de filas por tabla (fuente .207 - destino raw), 30d"
      description: "Positivo = destino atrasado respecto a .207. Cero es lo esperado la mayoria de los dias."
    plugin:
      kind: "TimeSeriesChart"
      spec:
        visual:
          display: "line"
    queries:
      - kind: "TimeSeriesQuery"
        spec:
          plugin:
            kind: "LokiTimeSeriesQuery"
            spec:
              query: 'avg by (dataset) (avg_over_time({job="soda_real"} | json | unwrap diff [1d]))'

rowDiffByDataset es el panel más útil de los dos: | json | unwrap diff extrae el campo numérico diff de cada línea de log (no solo cuenta líneas, como hacen los paneles basados en count_over_time) — muestra el hueco real en filas, no solo un semáforo pass/fail. Layout agregado en la grilla (y: 14 y y: 24, debajo de los 4 paneles originales que terminan en y: 4 + height: 10).

Verificación de que los datos llegaron

now=$(date +%s); start=$((now - 3600))
curl -s -G "http://127.0.0.1:3100/loki/api/v1/query_range" \
  --data-urlencode 'query={job="soda_real"}' \
  --data-urlencode "start=${start}000000000" --data-urlencode "end=${now}000000000" \
  --data-urlencode 'limit=20'

/loki/api/v1/query (instant query) rechaza consultas de log lines (“log queries are not supported as an instant query type”) — hace falta /query_range incluso para pedir solo las últimas líneas, gotcha de la API de Loki que no es específico de este check pero vale la pena anotar.

dvt-data-quality y data-health-status: quedan sin datos nuevos

Ambos seguían apuntando a job="dvt" — con dvt-checks/ decomisionado (cron quitado, imagen Docker borrada en una limpieza de disco anterior, ver DVT), esos dos dashboards van a mostrar “sin datos” para cualquier rango posterior al último job="dvt" real. No se tocaron ni se re-apuntaron a soda_real — decisión deliberada, para no mezclar la semántica de columna+schema (DVT) con la de solo-conteo (este check) bajo el mismo dashboard.

Véase también