Módulo 5: Point In Time Joins And Deduplication

De dónde vienen las filas duplicadas

Descripción

Las lecciones 2 a 4 de este módulo resolvieron un problema: cuál versión de la dimensión le corresponde a cada venta. Esta lección abre un problema distinto, que puede convivir perfectamente con un join punto-en-el-tiempo ya correcto: cuántas veces aparece la misma venta en los datos que llegan a Kiosko. Vas a construir un lote de órdenes donde tres de ellas llegaron dos veces —reenviadas, con una marca de tiempo de ingesta distinta cada vez— y vas a medir, con la misma consulta de grano del módulo 1, exactamente cuántas filas de más produce.

Conexión con el módulo. Esta lección es deliberadamente diagnóstica, no correctiva — construye la evidencia de que hay duplicados y explica de dónde salen, sin deduplicar todavía. La lección 6 toma exactamente el lote que esta lección construye y lo deduplica con ROW_NUMBER()/QUALIFY. Separar "detectar" de "corregir" en dos lecciones distintas es intencional: entender por qué existen los duplicados —antes de aprender a eliminarlos— evita tratar la deduplicación como un comando mágico que se aplica sin pensar en la causa.

Una analogía: el mismo paquete, tocado dos veces por error

Piensa en un mensajero que entrega un paquete, pero la app de rastreo no confirma la entrega a tiempo —quizás perdió señal, quizás el servidor tardó en responder—. El sistema, sin saber si el paquete realmente llegó, lo marca como "pendiente" y lo reintenta: el mensajero (u otro) toca la puerta una segunda vez, con el mismo paquete. Desde la perspectiva de quien recibe, no pasó nada raro —recibió un paquete—; pero desde la perspectiva del sistema de registro, ahora existen dos eventos de entrega para el mismo paquete, cuando en realidad ocurrió uno solo.

Eso es, con precisión, lo que le pasa a una orden que se reenvía: un sistema de ingesta que no confirma la recepción de forma confiable —la garantía técnica que en sistemas distribuidos se llama at-least-once delivery, "al menos una vez", en contraste con exactly-once, "exactamente una vez"— prefiere reenviar de más antes que perder datos. Es una decisión de diseño razonable del lado de la fuente: perder una orden real es mucho peor que procesarla dos veces. Pero eso significa que quien recibe esos datos —Kiosko, en este caso— tiene que estar preparado para encontrar la misma orden más de una vez, y saber qué hacer con eso.

Ejemplo trabajado: un lote con tres órdenes reenviadas

Construye un lote nuevo, raw_orders_batch, que simula cómo llegarían las cuarenta órdenes de la semana de Kiosko si el sistema de ingesta tuviera reintentos: cada orden llega una vez, con una marca ingested_at (un minuto después de order_ts, por regla fija) — excepto tres órdenes, que llegan dos veces, con un segundo ingested_at más tardío, simulando el reenvío.

# raw_orders_batch.py
from datetime import datetime, timedelta
import duckdb

from raw_orders import RAW_ORDERS

con = duckdb.connect()
con.execute("""
    CREATE TABLE raw_orders_batch (
        order_id VARCHAR, store_id VARCHAR, product_id VARCHAR,
        quantity INTEGER, unit_price DOUBLE, revenue DOUBLE,
        order_ts TIMESTAMP, ingested_at TIMESTAMP
    )
""")

# Cada orden llega una vez, con ingested_at = order_ts + 1 minuto (regla fija).
rows = []
for r in RAW_ORDERS:
    order_id, store_id, product_id, quantity, unit_price, order_ts_str = r
    order_ts = datetime.fromisoformat(order_ts_str)
    ingested_at = order_ts + timedelta(minutes=1)
    revenue = quantity * unit_price
    rows.append((order_id, store_id, product_id, quantity, unit_price, revenue, order_ts, ingested_at))

# Tres ordenes se REENVIAN -- mismo order_id + product_id, ingested_at posterior
# (retry / entrega "al menos una vez"), con un retraso fijo cada una.
RESEND_DELAY = {
    "ORD-1001": timedelta(minutes=35),
    "ORD-3001": timedelta(minutes=54),
    "ORD-6005": timedelta(minutes=49),
}
for r in RAW_ORDERS:
    order_id, store_id, product_id, quantity, unit_price, order_ts_str = r
    if order_id in RESEND_DELAY:
        order_ts = datetime.fromisoformat(order_ts_str)
        original_ingest = order_ts + timedelta(minutes=1)
        resend_ingest = original_ingest + RESEND_DELAY[order_id]
        revenue = quantity * unit_price
        rows.append((order_id, store_id, product_id, quantity, unit_price, revenue, order_ts, resend_ingest))

con.executemany("INSERT INTO raw_orders_batch VALUES (?, ?, ?, ?, ?, ?, ?, ?)", rows)

print("=== Las 3 llaves reenviadas, tal como llegan (dos filas cada una) ===")
print(con.sql("""
    SELECT order_id, product_id, ingested_at
    FROM raw_orders_batch
    WHERE order_id IN ('ORD-1001', 'ORD-3001', 'ORD-6005')
    ORDER BY order_id, ingested_at
"""))

Qué esperar. Al correr python3 raw_orders_batch.py, la salida es exactamente esta:

=== Las 3 llaves reenviadas, tal como llegan (dos filas cada una) ===
┌──────────┬────────────┬─────────────────────┐
│ order_id │ product_id │     ingested_at     │
│ varchar  │  varchar   │      timestamp      │
├──────────┼────────────┼─────────────────────┤
│ ORD-1001 │ P001       │ 2026-08-03 08:15:00 │
│ ORD-1001 │ P001       │ 2026-08-03 08:50:00 │
│ ORD-3001 │ P004       │ 2026-08-05 08:11:00 │
│ ORD-3001 │ P004       │ 2026-08-05 09:05:00 │
│ ORD-6005 │ P003       │ 2026-08-08 09:11:00 │
│ ORD-6005 │ P003       │ 2026-08-08 10:00:00 │
└──────────┴────────────┴─────────────────────┘

ORD-1001, ORD-3001 y ORD-6005 aparecen dos veces cada una, con la misma order_id y el mismo product_id, pero un ingested_at distinto — la marca de tiempo de cuándo el dato entró al sistema de Kiosko, que no es lo mismo que order_ts (cuándo ocurrió la venta). El resto de las órdenes de Kiosko —treinta y siete de las cuarenta— aparecen exactamente una vez, sin reenvío.

Diagrama: el grano roto, medido con la misma consulta del módulo 1

Ya conoces la consulta que declara el grano de una tabla de hechos —COUNT(*) contra COUNT(DISTINCT llave)— desde el módulo 1, lección 5. Aplícala aquí, con la misma llave compuesta (order_id + product_id) que usaste entonces:

print("\n=== Grano de raw_orders_batch: COUNT(*) vs COUNT(DISTINCT order_id-product_id) ===")
print(con.sql("""
    SELECT
        COUNT(*) AS total_rows,
        COUNT(DISTINCT order_id || '-' || product_id) AS distinct_order_product_lines
    FROM raw_orders_batch
"""))
=== Grano de raw_orders_batch: COUNT(*) vs COUNT(DISTINCT order_id-product_id) ===
┌────────────┬──────────────────────────────┐
│ total_rows │ distinct_order_product_lines │
│   int64    │            int64             │
├────────────┼──────────────────────────────┤
│         43 │                           40 │
└────────────┴──────────────────────────────┘

43 contra 40 — los dos números no coinciden, exactamente la señal que el módulo 1 enseñó a buscar. La diferencia, 3, es el número exacto de órdenes reenviadas. Identifica cuáles, con precisión, agrupando por la misma llave y filtrando las que aparecen más de una vez:

print("\n=== Identificando EXACTAMENTE cuales llaves estan duplicadas ===")
print(con.sql("""
    SELECT order_id || '-' || product_id AS order_product_key, COUNT(*) AS times_seen
    FROM raw_orders_batch
    GROUP BY order_product_key
    HAVING COUNT(*) > 1
    ORDER BY order_product_key
"""))
=== Identificando EXACTAMENTE cuales llaves estan duplicadas ===
┌────────────────────┬────────────┐
│ order_product_key  │ times_seen │
│      varchar       │   int64    │
├────────────────────┼────────────┤
│ ORD-1001-P001      │          2 │
│ ORD-3001-P004      │          2 │
│ ORD-6005-P003      │          2 │
└────────────────────┴────────────┘
flowchart TD
    A["Fuente: 40 ordenes reales de Kiosko"] --> B["Sistema de ingesta con reintentos\n(at-least-once delivery)"]
    B -->|"37 ordenes: 1 intento,\nconfirmacion a tiempo"| C["1 fila cada una"]
    B -->|"3 ordenes: la confirmacion\nno llego a tiempo, se reintenta"| D["2 filas cada una\n(mismo order_id+product_id,\ndistinto ingested_at)"]
    C --> E["raw_orders_batch: 43 filas totales\n40 llaves unicas, 3 duplicadas"]
    D --> E

Profundización: las causas reales de un lote duplicado, no solo esta

El escenario de esta lección —reintentos de un sistema de ingesta con entrega "al menos una vez"— es una causa real y común, pero no la única. Vale la pena nombrar las demás, porque la técnica de deduplicación de la lección 6 sirve para todas ellas por igual, sin importar cuál fue la causa específica:

  • Reintentos de red o de API (el caso de esta lección): un cliente HTTP no recibe confirmación a tiempo —por timeout, por un error transitorio del servidor— y reenvía la misma solicitud, sin saber si la primera llegó o no. La fuente prefiere el riesgo de duplicar sobre el riesgo de perder datos.
  • Múltiples corridas de extracción sobre la misma ventana de tiempo: un job de extracción que falla a la mitad, y se vuelve a correr desde el principio de la ventana en vez de reanudar donde se quedó, reintroduce las filas que ya había extraído antes de fallar.
  • Replay de un log de cambios (CDC): un sistema de Change Data Capture —mencionado en foundations y retomado a fondo en streaming-with-kafka-and-flink-guide— puede reproducir el mismo evento más de una vez si el consumidor se reinicia desde un punto de control (checkpoint) anterior al último evento realmente procesado.
  • Un mismo archivo cargado más de una vez por error humano o de automatización: alguien —o un cron mal configurado— vuelve a ejecutar la carga de orders_2026-08-03.csv sin darse cuenta de que ya se había cargado antes.

Todas estas causas comparten algo importante: ninguna corrompe el valor de las filas —la orden reenviada tiene, en todos los casos de esta lección, exactamente los mismos quantity, unit_price y revenue en ambas copias—. Esto las distingue de un problema distinto y más difícil: cuando un reenvío trae valores corregidos, no idénticos (por ejemplo, una cantidad distinta porque alguien corrigió un error de captura). Ese caso —una "corrección" disfrazada de duplicado— necesita una decisión de negocio explícita sobre cuál versión conservar, y queda fuera del alcance de esta lección: aquí, los duplicados son copias exactas, y la única decisión es cuál copia conservar, no cuál valor es el correcto.

Confirma esta propiedad —los duplicados son copias exactas, no correcciones— con una consulta:

print("\n=== Confirmando que las 2 copias de cada llave reenviada tienen el MISMO revenue ===")
print(con.sql("""
    SELECT order_id, product_id, COUNT(DISTINCT revenue) AS valores_distintos_de_revenue
    FROM raw_orders_batch
    WHERE order_id IN ('ORD-1001', 'ORD-3001', 'ORD-6005')
    GROUP BY order_id, product_id
    ORDER BY order_id
"""))
=== Confirmando que las 2 copias de cada llave reenviada tienen el MISMO revenue ===
┌──────────┬────────────┬───────────────────────────────┐
│ order_id │ product_id │ valores_distintos_de_revenue │
│ varchar  │  varchar   │             int64             │
├──────────┼────────────┼───────────────────────────────┤
│ ORD-1001 │ P001       │                             1 │
│ ORD-3001 │ P004       │                             1 │
│ ORD-6005 │ P003       │                             1 │
└──────────┴────────────┴───────────────────────────────┘

valores_distintos_de_revenue = 1 para las tres — cada par de filas duplicadas tiene exactamente el mismo revenue en ambas copias. Esta es la confirmación de que estás frente a un duplicado real, no una corrección: la deduplicación de la lección 6 es segura de aplicar precisamente porque no hay ninguna ambigüedad sobre qué valor es "el correcto" — ambas copias lo son, y sobra una.

Errores comunes

Deduplicar antes de confirmar que las filas son copias exactas. Qué pasa: alguien ve dos filas con el mismo order_id, asume que son duplicados en el sentido de esta lección, y las deduplica sin verificar si los valores realmente coinciden. Por qué pasa: dos filas con la misma llave "se ven" como duplicados a simple vista, y verificar los valores parece un paso extra innecesario. Cómo detectarlo: si tu proceso de deduplicación nunca comparó los valores de las filas "duplicadas" —solo sus llaves—, no puedes distinguir entre un reenvío real (mismo dato, dos veces) y una corrección legítima (mismo order_id, valores distintos) que debería tratarse de forma completamente diferente. Cómo corregirlo: antes de deduplicar cualquier lote, corre una consulta como la de esta lección —COUNT(DISTINCT <columna de valor>) agrupado por la llave sospechosa de duplicación— para confirmar que las copias son, de hecho, idénticas.

Usar order_id solo como la llave de grano, en vez de la llave compuesta. Qué pasa: alguien, al construir la consulta de grano de esta lección, usa COUNT(DISTINCT order_id) en vez de COUNT(DISTINCT order_id || '-' || product_id), repitiendo el mismo error que el módulo 1, lección 5, ya advirtió para fact_orders. Por qué pasa: con los datos actuales de Kiosko —donde cada orden tiene un solo producto— ambas consultas dan el mismo resultado (40 distintos), así que el error no se manifiesta con este dataset. Cómo detectarlo: si tu declaración del grano de raw_orders_batch no menciona product_id, corriste una verificación más débil de la que la estructura real de la tabla permite. Cómo corregirlo: usa siempre la llave compuesta más fina que el esquema soporte, exactamente como en el módulo 1 — la disciplina no cambia solo porque el contexto (deduplicación, en vez de declaración de grano) sí cambió.

Ignorar ingested_at y usar order_ts para decidir cuál copia es "la más reciente". Qué pasa: alguien, al preparar el terreno para deduplicar en la lección 6, planea usar order_ts —la fecha de la venta— para decidir qué copia conservar, en vez de ingested_at —la fecha de llegada del dato—. Por qué pasa: order_ts es la columna que ya conoces desde el módulo 1, y ingested_at es nueva en esta lección, así que es fácil recurrir a la que ya es familiar. Cómo detectarlo: en esta lección, las dos copias de cada orden reenviada tienen exactamente el mismo order_ts —la venta ocurrió una sola vez, en un solo instante—; lo único que distingue una copia de la otra es ingested_at. Si tu criterio de deduplicación usa order_ts, no tienes ninguna forma de distinguir entre las dos filas, porque son idénticas en esa columna. Cómo corregirlo: cuando deduplicas por reenvío, el criterio de "cuál copia conservar" debe basarse en una columna que distinga las copias entre sí — típicamente la marca de tiempo de ingesta, exactamente ingested_at en este caso. La lección 6 construye sobre esta distinción.

Ejercicios

Ejercicio 1 — Confirma que las 37 órdenes sin reenvío tienen exactamente una fila cada una. Usando raw_orders_batch, escribe una consulta que confirme que ninguna orden fuera de las tres reenviadas aparece más de una vez.

Ver solución
print(con.sql("""
    SELECT order_id || '-' || product_id AS order_product_key, COUNT(*) AS times_seen
    FROM raw_orders_batch
    WHERE order_id NOT IN ('ORD-1001', 'ORD-3001', 'ORD-6005')
    GROUP BY order_product_key
    HAVING COUNT(*) != 1
"""))

Salida esperada:

┌────────────────────┬────────────┐
│ order_product_key  │ times_seen │
│      varchar       │   int64    │
├────────────────────┼────────────┤
└────────────────────┴────────────┘
                0 rows

Cero filas — ninguna orden fuera de las tres reenviadas aparece un número de veces distinto a uno. Esto confirma, con evidencia negativa (la ausencia de resultados), que el problema de esta lección está acotado exactamente a las tres órdenes que el script marcó explícitamente para reenvío, no a un problema más amplio del lote completo.

Ejercicio 2 — Calcula cuánto se infla el revenue del lote si se suma sin deduplicar. Usando raw_orders_batch completo (43 filas), calcula SUM(revenue) y compáralo contra el revenue real de las cuarenta órdenes (106.15).

Ver solución
print(con.sql("SELECT ROUND(SUM(revenue), 2) AS total_revenue_con_duplicados FROM raw_orders_batch"))

Salida esperada:

┌──────────────────────────────┐
│ total_revenue_con_duplicados │
│            double            │
├──────────────────────────────┤
│                        119.05│
└──────────────────────────────┘

119.05 en vez de 106.15 — una inflación de 12.90, que es exactamente la suma del revenue de las tres órdenes reenviadas contada una vez de más (ORD-1001: 3 x 0.55 = 1.65; ORD-3001: 2 x 4.50 = 9.00; ORD-6005: 3 x 0.75 = 2.25; total 12.90). Este es el mismo tipo de error silencioso que ya viste con el fan-out de la lección 2 — un SUM sobre datos con duplicados produce un número perfectamente creíble, que está mal.

Ejercicio 3 — Explica por qué esta lección no deduplicó todavía, aunque ya tiene toda la evidencia para hacerlo. En 2-3 frases, explica la decisión pedagógica de separar "detectar duplicados" (esta lección) de "eliminarlos" (la lección 6).

Ver solución

Deduplicar sin haber confirmado primero por qué existen los duplicados —si son copias exactas de un reenvío, o correcciones legítimas con valores distintos— es aplicar una técnica sin entender si es la correcta para el problema específico. Esta lección construyó la evidencia completa: cuántas filas de más hay (43 contra 40), cuáles llaves están duplicadas exactamente, y que las copias son idénticas en valor (no correcciones). Con esa evidencia ya confirmada, la lección 6 puede aplicar ROW_NUMBER()/QUALIFY con la confianza de que eliminar las copias de más no descarta ninguna información real — una decisión que solo tiene sentido después de haber hecho el diagnóstico, no antes.

Resumen y siguiente paso

Esta lección construyó un lote de órdenes con tres reenvíos reales —ORD-1001, ORD-3001, ORD-6005, cada una con dos copias idénticas en valor pero distinta ingested_at— y confirmó, con la misma consulta de grano del módulo 1, que el lote tiene 43 filas para 40 combinaciones únicas de order_id+product_id. Nombraste las causas más comunes de este tipo de duplicación —reintentos de red, extracciones repetidas, replay de CDC, cargas duplicadas por error humano— y confirmaste que, en este caso, las copias son exactas: el problema es "cuál conservar", no "cuál es correcta".

Antes de avanzar deberías poder: aplicar la consulta de grano del módulo 1 a cualquier lote nuevo para detectar duplicados; nombrar al menos tres causas reales de duplicación en un pipeline de datos; y explicar la diferencia entre un duplicado exacto (esta lección) y una corrección disfrazada de duplicado (fuera de alcance).

La lección 6 toma este mismo lote, raw_orders_batch, y lo deduplica de verdad: ROW_NUMBER() OVER (PARTITION BY order_id, product_id ORDER BY ingested_at DESC) combinado con QUALIFY, conservando la copia más reciente de cada reenvío y restaurando el grano exacto de cuarenta filas.

Recursos