Módulo 5: Point In Time Joins And Deduplication
Deduplicando con ROW_NUMBER y QUALIFY
Descripción
La lección 5 diagnosticó el problema con evidencia completa: raw_orders_batch tiene 43 filas para 40 combinaciones únicas de order_id+product_id, con tres órdenes reenviadas, cada una con dos copias idénticas en valor. Esta lección lo resuelve: ROW_NUMBER(), una función de ventana que numera las filas dentro de cada grupo, combinada con QUALIFY, la cláusula de DuckDB que filtra sobre el resultado de una función de ventana sin necesitar una subconsulta. El resultado —verificado, no asumido— es que el lote vuelve a tener exactamente cuarenta filas, con el revenue correcto.
Conexión con el módulo. Esta lección resuelve la segunda de las dos filas del checklist que le corresponden a este módulo —"deduplicación explícita de filas repetidas"—, la primera fue el join punto-en-el-tiempo de las lecciones 2 a 4. Las dos técnicas no dependen una de la otra —puedes deduplicar sin nunca haber unido contra una dimensión historizada—, pero conviven en el mismo módulo porque ambas resuelven el mismo tipo de pregunta: "¿qué fila, exactamente, debería representar este evento?", cuando hay más de una candidata.
Una analogía: quedarte con la última actualización, no con todas
Piensa en un documento compartido que varias personas editan a la vez, y cada vez que alguien guarda, el sistema crea una nueva versión con marca de tiempo —sin borrar las anteriores—. Si alguien te pide "el documento", nadie espera que le entregues las quince versiones guardadas ese día: esperan la más reciente, la que refleja el estado final después de todas las ediciones. Las versiones anteriores no están "mal" —cada una fue, en su momento, la versión correcta—, pero solo una de ellas es la que importa para responder "¿cómo quedó el documento?" hoy.
ROW_NUMBER() OVER (PARTITION BY order_id, product_id ORDER BY ingested_at DESC) hace exactamente eso con las órdenes reenviadas de Kiosko: agrupa (PARTITION BY) todas las copias de la misma orden, las numera de la más reciente a la más antigua (ORDER BY ingested_at DESC, así que la más reciente recibe el número 1), y QUALIFY ... = 1 se queda solo con esa — la versión final, descartando las copias intermedias sin perder ninguna información que la versión final no tenga ya.
Ejemplo trabajado: ROW_NUMBER, visto fila por fila, antes de filtrar
Reconstruye raw_orders_batch exactamente como en la lección 5 —el mismo script, raw_orders_batch.py—. Antes de filtrar nada, mira lo que ROW_NUMBER() calcula para las tres llaves reenviadas, sin QUALIFY todavía, para entender exactamente qué numera:
# dedup_qualify.py
# (raw_orders_batch ya construido como en la leccion 5)
print("=== ROW_NUMBER() sobre las 3 llaves reenviadas, SIN filtrar todavia ===")
print(con.sql("""
SELECT order_id, product_id, ingested_at,
ROW_NUMBER() OVER (PARTITION BY order_id, product_id ORDER BY ingested_at DESC) AS rn
FROM raw_orders_batch
WHERE order_id IN ('ORD-1001', 'ORD-3001', 'ORD-6005')
ORDER BY order_id, rn
"""))
=== ROW_NUMBER() sobre las 3 llaves reenviadas, SIN filtrar todavia ===
┌──────────┬────────────┬─────────────────────┬───────┐
│ order_id │ product_id │ ingested_at │ rn │
│ varchar │ varchar │ timestamp │ int64 │
├──────────┼────────────┼─────────────────────┼───────┤
│ ORD-1001 │ P001 │ 2026-08-03 08:50:00 │ 1 │
│ ORD-1001 │ P001 │ 2026-08-03 08:15:00 │ 2 │
│ ORD-3001 │ P004 │ 2026-08-05 09:05:00 │ 1 │
│ ORD-3001 │ P004 │ 2026-08-05 08:11:00 │ 2 │
│ ORD-6005 │ P003 │ 2026-08-08 10:00:00 │ 1 │
│ ORD-6005 │ P003 │ 2026-08-08 09:11:00 │ 2 │
└──────────┴────────────┴─────────────────────┴───────┘
Fíjate en el orden: para cada order_id+product_id, la fila con el ingested_at más tardío —la copia reenviada, la más reciente en llegar— recibe rn = 1; la copia original recibe rn = 2. Eso es exactamente lo que ORDER BY ingested_at DESC (descendente) produce: la fecha más grande primero, así que rn = 1 siempre corresponde a "la última vez que este dato llegó".
Ahora, QUALIFY: filtra el resultado de la función de ventana, quedándose solo con rn = 1, sin necesitar una subconsulta o un WITH para lograrlo.
print("\n=== QUALIFY aplicado -- vuelve a 40 filas ===")
print(con.sql("""
WITH deduped AS (
SELECT *
FROM raw_orders_batch
QUALIFY ROW_NUMBER() OVER (PARTITION BY order_id, product_id ORDER BY ingested_at DESC) = 1
)
SELECT COUNT(*) AS total_rows FROM deduped
"""))
=== QUALIFY aplicado -- vuelve a 40 filas ===
┌────────────┐
│ total_rows │
│ int64 │
├────────────┤
│ 40 │
└────────────┘
Cuarenta filas — exactamente el número de líneas de orden únicas que raw_orders_batch siempre tuvo, según la lección 5. Verifica el grano de nuevo, con la misma consulta de siempre, ahora sobre el resultado deduplicado, y confirma el revenue correcto:
print("\n=== Grano verificado de nuevo, post-dedup ===")
print(con.sql("""
WITH deduped AS (
SELECT *
FROM raw_orders_batch
QUALIFY ROW_NUMBER() OVER (PARTITION BY order_id, product_id ORDER BY ingested_at DESC) = 1
)
SELECT
COUNT(*) AS total_rows,
COUNT(DISTINCT order_id || '-' || product_id) AS distinct_order_product_lines,
ROUND(SUM(revenue), 2) AS total_revenue
FROM deduped
"""))
=== Grano verificado de nuevo, post-dedup ===
┌────────────┬──────────────────────────────┬────────────────┐
│ total_rows │ distinct_order_product_lines │ total_revenue │
│ int64 │ int64 │ double │
├────────────┼──────────────────────────────┼────────────────┤
│ 40 │ 40 │ 106.15 │
└────────────┴──────────────────────────────┴────────────────┘
total_rows y distinct_order_product_lines vuelven a coincidir —el grano está restaurado—, y total_revenue vuelve a ser 106.15, el número de referencia de siempre, no los 119.05 inflados que la lección 5 midió sobre el lote sin deduplicar.
Diagrama: qué hace cada pieza de la consulta
flowchart TD
A["raw_orders_batch: 43 filas"] --> B["PARTITION BY order_id, product_id\nagrupa las copias de la misma linea de orden"]
B --> C["ORDER BY ingested_at DESC\ndentro de cada grupo, la mas reciente primero"]
C --> D["ROW_NUMBER()\nnumera 1, 2, 3... dentro de cada grupo"]
D --> E{"QUALIFY rn = 1"}
E -->|"rn = 1\n(la copia mas reciente)"| F["Se conserva: 40 filas"]
E -->|"rn > 1\n(copias mas antiguas)"| G["Se descarta: 3 filas"]
Profundización: por qué QUALIFY, y qué reemplaza exactamente
Antes de que QUALIFY existiera como cláusula, filtrar por el resultado de una función de ventana requería envolver la consulta en una subconsulta o un WITH —exactamente lo que el CTE deduped de esta lección hace por dentro—, porque las funciones de ventana, a diferencia de WHERE, se calculan después de que las filas ya fueron seleccionadas, no antes. QUALIFY existe, en palabras de la documentación oficial de DuckDB, con la misma relación que HAVING tiene con GROUP BY: así como HAVING filtra el resultado de una función agregada sin necesitar una subconsulta, QUALIFY filtra el resultado de una función de ventana sin necesitarla tampoco. La posición exacta de QUALIFY dentro de un SELECT completo, según esa misma documentación, es después de WINDOW y antes de ORDER BY:
SELECT select_list FROM tables WHERE condition GROUP BY groups
HAVING group_filter WINDOW window_expression QUALIFY qualify_filter
ORDER BY order_expression LIMIT n
Esta es la forma equivalente, sin QUALIFY, que habrías tenido que escribir en un motor que no lo soporte:
-- Equivalente SIN QUALIFY, con una subconsulta explicita
SELECT order_id, store_id, product_id, quantity, unit_price, revenue, order_ts, ingested_at
FROM (
SELECT *, ROW_NUMBER() OVER (PARTITION BY order_id, product_id ORDER BY ingested_at DESC) AS rn
FROM raw_orders_batch
) t
WHERE rn = 1
Las dos formas producen exactamente el mismo resultado — QUALIFY no cambia la lógica, solo evita tener que nombrar y envolver una subconsulta únicamente para poder filtrar por rn. Vale la pena notar una restricción de DuckDB que aparece si intentas combinar QUALIFY con una función de agregación sin GROUP BY en la misma consulta —por ejemplo, SELECT COUNT(*) FROM t QUALIFY ROW_NUMBER() ... = 1, sin un CTE de por medio—: DuckDB pide que las columnas usadas dentro de la función de ventana también aparezcan en un GROUP BY, porque trata la consulta completa como una agregación. La forma correcta, la que usa esta lección, evita ese conflicto separando el filtrado (QUALIFY, dentro del CTE deduped) de la agregación final (COUNT(*), fuera de él) — dos pasos distintos, cada uno con su propia responsabilidad.
Por qué ORDER BY ingested_at DESC, y no ASC
Fíjate en la dirección del ORDER BY: DESC, descendente, no ASC. La elección no es arbitraria — depende de qué copia representa "la verdad" cuando hay más de una. En el escenario de esta lección, las copias son idénticas en valor (la lección 5 lo confirmó), así que cualquiera de las dos sería técnicamente correcta de conservar. Pero la convención de "quedarse con el registro más reciente" —ORDER BY <marca_de_tiempo> DESC, rn = 1— es la que generaliza correctamente al caso más común en producción, donde un reenvío puede traer una corrección real (no solo una copia exacta): si un sistema fuente corrige un error y reenvía la misma orden con un valor distinto, quedarse con la copia más reciente es, casi siempre, la decisión correcta — es la versión más actualizada de la verdad, según la fuente. Usar ASC en su lugar conservaría deliberadamente la copia más antigua, una decisión válida solo si tu regla de negocio es "el primer registro que llega es el que cuenta, sin importar reenvíos posteriores" — una regla distinta, con sus propios casos de uso, pero no la que esta lección adopta.
Errores comunes
Usar PARTITION BY order_id solo, sin product_id. Qué pasa: alguien escribe ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY ingested_at DESC), sin incluir product_id en el PARTITION BY, repitiendo el mismo error de llave incompleta que el módulo 1 y la lección 5 de este módulo ya advirtieron para el grano de fact_orders. Por qué pasa: con los datos actuales de Kiosko —donde cada orden tiene un solo producto—, order_id solo produce el mismo resultado que la llave compuesta, así que el error es invisible con este dataset. Cómo detectarlo: si tu PARTITION BY no incluye todas las columnas que forman el grano real de la tabla (aquí, order_id y product_id), y algún día una orden de Kiosko tuviera más de un producto, el ROW_NUMBER() agruparía líneas de orden completamente distintas como si fueran duplicados de la misma cosa —descartando, por error, una línea de producto real, no un duplicado—. Cómo corregirlo: el PARTITION BY de una deduplicación debe usar exactamente la misma llave que la consulta de grano usa para COUNT(DISTINCT ...) — en Kiosko, siempre order_id y product_id juntos, nunca uno solo.
Olvidar el ORDER BY dentro de ROW_NUMBER(), o usar una columna que no distingue las copias. Qué pasa: alguien escribe ROW_NUMBER() OVER (PARTITION BY order_id, product_id ORDER BY order_ts) —usando order_ts en vez de ingested_at—, tal como advirtió la lección 5. Como las dos copias de cada orden reenviada comparten el mismo order_ts (la venta ocurrió una sola vez), DuckDB numera las filas con ese empate en un orden que no está garantizado de ser determinista sin una columna de desempate. Por qué pasa: order_ts es la columna más familiar, y no es obvio, sin pensarlo, que las dos copias comparten ese valor exactamente. Cómo detectarlo: si corres la misma consulta de deduplicación dos veces y el resultado (cuál fila específica queda con rn = 1) no es siempre el mismo, tu ORDER BY tiene un empate sin desempatar. Cómo corregirlo: el ORDER BY de un ROW_NUMBER() de deduplicación debe usar una columna que sí distinga las copias entre sí de forma única — ingested_at, en el caso de esta lección, porque cada reenvío llegó en un instante distinto.
Asumir que QUALIFY reemplaza a WHERE, en vez de complementarlo. Qué pasa: alguien intenta escribir WHERE ROW_NUMBER() OVER (...) = 1 directamente, sin QUALIFY, y DuckDB rechaza la consulta con un error de sintaxis. Por qué pasa: WHERE filtra filas antes de que las funciones de ventana se calculen, así que en el momento en que WHERE se evalúa, ROW_NUMBER() todavía no existe como valor disponible — es imposible filtrar por algo que todavía no se calculó. Cómo detectarlo: si DuckDB reporta un error de "column not found" o de sintaxis al intentar usar una función de ventana dentro de un WHERE, ese es exactamente el síntoma. Cómo corregirlo: cualquier filtro que dependa del resultado de una función de ventana —ROW_NUMBER(), RANK(), o cualquier otra— necesita QUALIFY (o la subconsulta/WITH equivalente), nunca WHERE directamente.
Ejercicios
Ejercicio 1 — Confirma que la copia conservada de cada llave reenviada es la de ingested_at más reciente. Usando el resultado deduplicado, escribe una consulta que muestre, para las tres llaves reenviadas, el ingested_at que quedó después de aplicar QUALIFY.
Ver solución
print(con.sql("""
SELECT order_id, ingested_at
FROM raw_orders_batch
QUALIFY ROW_NUMBER() OVER (PARTITION BY order_id, product_id ORDER BY ingested_at DESC) = 1
AND order_id IN ('ORD-1001', 'ORD-3001', 'ORD-6005')
ORDER BY order_id
"""))
Salida esperada:
┌──────────┬─────────────────────┐
│ order_id │ ingested_at │
│ varchar │ timestamp │
├──────────┼─────────────────────┤
│ ORD-1001 │ 2026-08-03 08:50:00 │
│ ORD-3001 │ 2026-08-05 09:05:00 │
│ ORD-6005 │ 2026-08-08 10:00:00 │
└──────────┴─────────────────────┘
Los tres ingested_at son, exactamente, los más tardíos de cada par —comparando contra la tabla de la lección 5 (08:15:00/08:50:00 para ORD-1001, 08:11:00/09:05:00 para ORD-3001, 09:11:00/10:00:00 para ORD-6005)—, confirmando que ORDER BY ingested_at DESC + rn = 1 conserva, sin excepción, la copia más reciente de cada reenvío.
Ejercicio 2 — Deduplica usando ASC en vez de DESC, y explica la diferencia. Corre la misma consulta de deduplicación, pero con ORDER BY ingested_at ASC (ascendente), y confirma qué ingested_at queda conservado para las tres llaves reenviadas.
Ver solución
print(con.sql("""
SELECT order_id, ingested_at
FROM raw_orders_batch
QUALIFY ROW_NUMBER() OVER (PARTITION BY order_id, product_id ORDER BY ingested_at ASC) = 1
AND order_id IN ('ORD-1001', 'ORD-3001', 'ORD-6005')
ORDER BY order_id
"""))
Salida esperada:
┌──────────┬─────────────────────┐
│ order_id │ ingested_at │
│ varchar │ timestamp │
├──────────┼─────────────────────┤
│ ORD-1001 │ 2026-08-03 08:15:00 │
│ ORD-3001 │ 2026-08-05 08:11:00 │
│ ORD-6005 │ 2026-08-08 09:11:00 │
└──────────┴─────────────────────┘
Con ASC, se conserva la copia más antigua de cada reenvío —los mismos tres ingested_at de la lección 5, pero los primeros, no los segundos—. En este caso específico, quantity, unit_price y revenue son idénticos entre ambas copias (confirmado en la lección 5), así que el resultado final del reporte —el revenue total, 106.15— sería idéntico con ASC o con DESC. La diferencia solo importaría si las copias tuvieran valores distintos entre sí, el escenario de "corrección disfrazada de duplicado" que la lección 5 nombró como fuera de alcance.
Ejercicio 3 — Explica por qué total_rows = 40 después de QUALIFY no es, por sí solo, prueba suficiente de que la deduplicación fue correcta. En 2-3 frases, describe qué otra verificación necesitarías correr para confirmar que las filas conservadas son las correctas, no solo que el conteo es correcto.
Ver solución
Que total_rows vuelva a ser 40 solo prueba que la cantidad de filas es correcta — no prueba que se haya conservado la copia correcta de cada reenvío, ni que no se haya descartado, por error, una línea de orden que en realidad era única (no un duplicado). Para confirmar eso, hace falta una verificación adicional sobre el contenido: comparar el revenue total post-deduplicación contra el valor de referencia conocido (106.15, exactamente lo que esta lección hizo), y, más específico todavía, confirmar que las tres llaves reenviadas conservaron el ingested_at esperado según el criterio de orden elegido (DESC), como hizo el ejercicio 1. Un conteo correcto con contenido equivocado sigue siendo un resultado incorrecto.
Resumen y siguiente paso
Esta lección deduplicó el lote de la lección 5 con ROW_NUMBER() OVER (PARTITION BY order_id, product_id ORDER BY ingested_at DESC) combinado con QUALIFY ... = 1, y confirmó, con la misma consulta de grano de siempre, que el resultado tiene exactamente cuarenta filas —el número correcto—, con el revenue total de vuelta en 106.15. Viste, además, por qué QUALIFY existe (evitar la subconsulta que un filtro sobre una función de ventana necesitaría de otra forma) y por qué la dirección del ORDER BY (DESC, "el más reciente gana") es una decisión de negocio, no un detalle técnico arbitrario.
Antes de avanzar deberías poder: escribir de memoria la estructura ROW_NUMBER() OVER (PARTITION BY ... ORDER BY ...) ... QUALIFY = 1; explicar por qué PARTITION BY debe usar la misma llave compuesta que la consulta de grano; y explicar la diferencia entre ORDER BY ... DESC (conserva lo más reciente) y ASC (conserva lo más antiguo) en el contexto de deduplicación.
La lección 7 introduce una herramienta relacionada, pero distinta: ANTI JOIN/SEMI JOIN, dos tipos de JOIN nativos de DuckDB que sirven tanto para encontrar filas de hechos sin ninguna versión de dimensión que las cubra (retomando el problema de la lección 4) como para detectar, con precisión, qué filas de un lote nuevo representan un cambio real frente a lo que ya está cargado.
Recursos
- DuckDB — documentación oficial de la cláusula
QUALIFY— la referencia exacta de sintaxis y posición dentro de unSELECT, incluyendo la comparación explícita conHAVING. duckdb.org/docs/current/sql/query_syntax/qualify. En inglés. - DuckDB — documentación de funciones de ventana, incluida
ROW_NUMBER()— la referencia completa de sintaxis dePARTITION BYyORDER BYdentro deOVER (...). duckdb.org/docs/current/sql/functions/window_functions. En inglés. - DuckDB — documentación oficial del cliente Python, usada para construir y deduplicar
raw_orders_batchen esta lección. duckdb.org/docs/current/clients/python/overview. En inglés.