Módulo 6: Accumulating And Cumulative Patterns
Actualizando milestones in-place
Descripción
La lección 3 construyó fact_sessions completo con una sola consulta agregada — el equivalente de un recálculo total, como si Kiosko procesara toda la semana de eventos de una sola vez, al final. En producción, los eventos no llegan así: llegan uno a la vez, en el orden en que ocurren, y el pipeline que mantiene fact_sessions los procesa a medida que llegan, sin esperar a tener la semana completa. Esta lección reconstruye la misma tabla, evento por evento, en orden cronológico real — y verifica, con una comparación directa, que produce exactamente el mismo resultado que el recálculo agregado de la lección 3.
Conexión con el módulo. Esta lección demuestra que el mecanismo mínimo de la lección 2 (INSERT en el primer milestone, UPDATE en los siguientes) escala, sin ningún cambio de lógica, a los 32 eventos completos de Kiosko — y que el resultado es idéntico, fila por fila, al que produjo el enfoque agregado de la lección 3. Esta equivalencia es la base de la lección 5, que aplica la misma idea de "procesar incrementalmente" a un patrón distinto.
Una analogía: la lista de casilleros, revisada uno a la vez
Imagina que Kiosko tuviera un tablero físico con un casillero por sesión activa, y un empleado revisando la fila de eventos que llegan del sistema, uno tras otro, en el orden en que ocurrieron. Cuando llega el primer evento de una sesión nueva —siempre un page_view—, el empleado abre un casillero nuevo con esa sesión. Cuando llega un evento de una sesión que ya tiene casillero, el empleado no abre uno nuevo: escribe el dato encima del casillero existente. El empleado nunca necesita saber, de antemano, cuántos eventos va a recibir en total, ni en qué orden exacto van a llegar los milestones de cada sesión — solo necesita saber, para cada evento que llega, si la sesión ya tiene casillero o no.
Eso es, exactamente, lo que esta lección construye en código: un bucle que revisa cada evento una vez, decide si la sesión ya existe en fact_sessions, y actúa en consecuencia — INSERT si no existe, UPDATE si ya existe.
Ejemplo trabajado: procesando 32 eventos, uno a la vez, en orden cronológico
El algoritmo central de esta lección tiene una regla simple: por cada evento, si la sesión no existe todavía en fact_sessions, se inserta una fila nueva (siempre con un page_view, el primer milestone del funnel); si la sesión ya existe, se actualiza la columna correspondiente al tipo de evento, sin tocar el resto de la fila.
# fact_sessions_incremental.py -- continua sobre RAW_EVENTS y store_for_session (leccion 3)
from datetime import datetime
import duckdb
con = duckdb.connect()
con.execute("""
CREATE TABLE fact_sessions (
session_id VARCHAR PRIMARY KEY,
store_id VARCHAR,
session_date DATE,
view_ts TIMESTAMP,
add_to_cart_ts TIMESTAMP,
purchase_ts TIMESTAMP,
is_converted BOOLEAN
)
""")
# Orden CRONOLOGICO -- exactamente como los eventos llegarian en produccion,
# no el orden en que estan declarados en events.py
events_sorted = sorted(RAW_EVENTS, key=lambda r: r[3])
inserts, updates = 0, 0
for event_id, event_type, session_id, event_ts_str in events_sorted:
event_ts = datetime.fromisoformat(event_ts_str)
exists = con.execute(
"SELECT 1 FROM fact_sessions WHERE session_id = ?", [session_id]
).fetchone()
if exists is None:
# Primer evento de esta sesion -- SIEMPRE page_view en este dataset -- INSERT
con.execute(
"INSERT INTO fact_sessions VALUES (?, ?, ?, ?, NULL, NULL, false)",
[session_id, store_for_session(session_id), event_ts.date(), event_ts],
)
inserts += 1
elif event_type == "add_to_cart":
# La sesion ya existe -- UPDATE in place, nunca una fila nueva
con.execute(
"UPDATE fact_sessions SET add_to_cart_ts = ? WHERE session_id = ?",
[event_ts, session_id],
)
updates += 1
elif event_type == "purchase":
con.execute(
"UPDATE fact_sessions SET purchase_ts = ?, is_converted = true WHERE session_id = ?",
[event_ts, session_id],
)
updates += 1
print(f"eventos procesados: {len(events_sorted)} (INSERT: {inserts}, UPDATE: {updates})")
print(f"filas finales en fact_sessions: {con.sql('SELECT COUNT(*) FROM fact_sessions').fetchone()[0]}\n")
con.sql("SELECT * FROM fact_sessions ORDER BY session_id").show(max_width=300)
Qué esperar.
eventos procesados: 32 (INSERT: 17, UPDATE: 15)
filas finales en fact_sessions: 17
┌────────────┬──────────┬──────────────┬─────────────────────┬─────────────────────┬─────────────────────┬──────────────┐
│ session_id │ store_id │ session_date │ view_ts │ add_to_cart_ts │ purchase_ts │ is_converted │
│ varchar │ varchar │ date │ timestamp │ timestamp │ timestamp │ boolean │
├────────────┼──────────┼──────────────┼─────────────────────┼─────────────────────┼─────────────────────┼──────────────┤
│ SESS-01 │ S01 │ 2026-08-03 │ 2026-08-03 08:00:12 │ 2026-08-03 08:02:45 │ 2026-08-03 08:03:10 │ true │
│ SESS-02 │ S02 │ 2026-08-03 │ 2026-08-03 08:05:00 │ NULL │ NULL │ false │
│ SESS-03 │ S03 │ 2026-08-04 │ 2026-08-04 08:10:00 │ 2026-08-04 08:12:30 │ 2026-08-04 08:13:05 │ true │
│ SESS-04 │ S01 │ 2026-08-04 │ 2026-08-04 08:20:00 │ NULL │ NULL │ false │
│ SESS-05 │ S02 │ 2026-08-04 │ 2026-08-04 08:45:00 │ NULL │ NULL │ false │
│ SESS-06 │ S03 │ 2026-08-05 │ 2026-08-05 08:00:00 │ NULL │ NULL │ false │
│ SESS-07 │ S01 │ 2026-08-05 │ 2026-08-05 08:15:00 │ 2026-08-05 08:16:20 │ NULL │ false │
│ SESS-08 │ S02 │ 2026-08-06 │ 2026-08-06 08:05:00 │ 2026-08-06 08:07:15 │ 2026-08-06 08:08:00 │ true │
│ SESS-09 │ S03 │ 2026-08-06 │ 2026-08-06 08:30:00 │ NULL │ NULL │ false │
│ SESS-10 │ S01 │ 2026-08-07 │ 2026-08-07 08:00:00 │ 2026-08-07 08:03:10 │ 2026-08-07 08:04:00 │ true │
│ SESS-11 │ S02 │ 2026-08-07 │ 2026-08-07 08:20:00 │ 2026-08-07 08:22:00 │ NULL │ false │
│ SESS-12 │ S03 │ 2026-08-07 │ 2026-08-07 08:50:00 │ NULL │ NULL │ false │
│ SESS-13 │ S01 │ 2026-08-08 │ 2026-08-08 07:55:00 │ 2026-08-08 07:58:00 │ 2026-08-08 07:59:10 │ true │
│ SESS-14 │ S02 │ 2026-08-08 │ 2026-08-08 08:10:00 │ 2026-08-08 08:12:45 │ 2026-08-08 08:13:30 │ true │
│ SESS-15 │ S03 │ 2026-08-08 │ 2026-08-08 08:40:00 │ NULL │ NULL │ false │
│ SESS-16 │ S01 │ 2026-08-09 │ 2026-08-09 09:00:00 │ NULL │ NULL │ false │
│ SESS-17 │ S02 │ 2026-08-09 │ 2026-08-09 09:20:00 │ 2026-08-09 09:22:00 │ NULL │ false │
└────────────┴──────────┴──────────────┴─────────────────────┴─────────────────────┴─────────────────────┴──────────────┘
Detente en la primera línea: INSERT: 17, UPDATE: 15. Diecisiete sesiones nuevas —una por cada page_view que abre una sesión— y quince actualizaciones —los 9 add_to_cart más los 6 purchase que le siguieron a alguna de esas sesiones—. 17 + 15 = 32, el total exacto de eventos, sin que ninguno se haya perdido ni contado dos veces. Y el resultado final —diecisiete filas, con exactamente los mismos valores en cada columna— es idéntico a la tabla que la lección 3 construyó con MAX(CASE WHEN ...).
Verificando la equivalencia entre los dos enfoques
Que ambas tablas "se vean iguales" no es suficiente evidencia — esta lección lo confirma comparándolas fila por fila:
# fact_sessions_batch.py -- reconstruye la version de la leccion 3, con otro nombre de tabla
con.execute("""
CREATE TABLE fact_sessions_batch AS
SELECT
e.session_id, m.store_id, MIN(CAST(e.event_ts AS DATE)) AS session_date,
MAX(CASE WHEN e.event_type = 'page_view' THEN e.event_ts END) AS view_ts,
MAX(CASE WHEN e.event_type = 'add_to_cart' THEN e.event_ts END) AS add_to_cart_ts,
MAX(CASE WHEN e.event_type = 'purchase' THEN e.event_ts END) AS purchase_ts,
MAX(CASE WHEN e.event_type = 'purchase' THEN true ELSE false END) AS is_converted
FROM events e JOIN session_store_map m ON e.session_id = m.session_id
GROUP BY e.session_id, m.store_id
""")
mismatches = con.sql("""
SELECT COUNT(*) FROM fact_sessions incr
JOIN fact_sessions_batch batch ON incr.session_id = batch.session_id
WHERE incr.store_id != batch.store_id
OR incr.view_ts != batch.view_ts
OR incr.add_to_cart_ts IS DISTINCT FROM batch.add_to_cart_ts
OR incr.purchase_ts IS DISTINCT FROM batch.purchase_ts
OR incr.is_converted != batch.is_converted
""").fetchone()[0]
print(f"filas donde la version incremental difiere de la version agregada: {mismatches}")
assert mismatches == 0, "las dos versiones de fact_sessions no coinciden"
print("Verificacion OK: el enfoque evento-por-evento produce EXACTAMENTE la misma tabla que MAX(CASE WHEN...)")
Qué esperar.
filas donde la version incremental difiere de la version agregada: 0
Verificacion OK: el enfoque evento-por-evento produce EXACTAMENTE la misma tabla que MAX(CASE WHEN...)
Cero diferencias. Esto confirma algo importante sobre el patrón accumulating snapshot: el orden mecánico en que construyes la tabla —todo de una vez con una agregación, o un evento a la vez con INSERT/UPDATE— no cambia el resultado, siempre que el conjunto de eventos de entrada sea el mismo. Lo que sí cambia es cuál de los dos enfoques se parece más a como funciona un sistema de producción real: un pipeline que recibe eventos continuamente (por ejemplo, uno que lee de una cola de mensajes o un CDC de eventos) usa el patrón INSERT/UPDATE de esta lección, procesando cada evento a medida que llega, sin esperar a tener todos los datos del día completos.
Diagrama: dos caminos, el mismo destino
flowchart TD
subgraph batch["Leccion 3: enfoque agregado"]
A["32 events completos\n(toda la semana disponible)"] --> B["MAX(CASE WHEN...)\nGROUP BY session_id"]
B --> C["fact_sessions_batch\n17 filas"]
end
subgraph incremental["Leccion 4: enfoque incremental"]
D["Evento 1 llega"] --> E{"la sesion ya existe?"}
E -->|"No"| F["INSERT\nfila nueva incompleta"]
E -->|"Si"| G["UPDATE\nla columna del milestone"]
F --> H["Siguiente evento..."]
G --> H
H -.->|"repite 32 veces"| D
H --> I["fact_sessions\n17 filas"]
end
C -.->|"0 diferencias, verificado"| I
Profundización: por qué el orden cronológico de los eventos importa aquí, y no en la lección 3
En la lección 3, el orden en que events aparece en la tabla no afecta el resultado — MAX(CASE WHEN ...) mira todos los eventos de una sesión a la vez, sin importar en qué secuencia los procesa el motor internamente. En esta lección, el orden sí importa, y por una razón concreta: el algoritmo decide si hacer INSERT o UPDATE mirando si la sesión ya existe en ese momento del procesamiento — una decisión que depende de qué eventos ya se procesaron antes. Si esta lección procesara los eventos fuera de orden cronológico (por ejemplo, un purchase antes que su page_view correspondiente), el algoritmo intentaría un UPDATE sobre una sesión que todavía no existe, y no encontraría ninguna fila que actualizar — el purchase se perdería silenciosamente, sin ningún error visible.
Esta es una limitación real, no un detalle académico: en un sistema de producción, los eventos pueden llegar fuera de orden por razones de red o de particionamiento de una cola de mensajes. La solución completa a ese problema —qué hacer cuando un milestone llega "antes" que el que debería precederlo— es exactamente el problema de dimensiones de llegada tardía que el módulo 5 ya nombró para dim_product_scd, aplicado aquí a hechos en vez de dimensiones; queda fuera del alcance de esta guía resolverlo a fondo (pertenece a streaming-with-kafka-and-flink-guide, donde el orden de llegada de eventos en tiempo real sí se maneja con garantías del sistema de streaming). Esta lección asume, como todo el hilo ejecutable de esta guía, que los eventos se procesan en el orden en que ocurrieron.
Errores comunes
Procesar events en el orden en que aparece en el archivo, sin ordenar por event_ts. Qué pasa: alguien itera sobre RAW_EVENTS directamente, en el orden en que está declarado en Python, sin aplicar sorted(..., key=lambda r: r[3]) primero. Por qué pasa: RAW_EVENTS ya está declarado, visualmente, en un orden que parece cronológico —agrupado por día—, así que parece innecesario ordenarlo de nuevo. Cómo detectarlo: si tu resultado difiere del de esta lección en alguna sesión, revisa si estás confiando en el orden de declaración en vez de ordenar explícitamente por event_ts — el orden de declaración en un archivo Python nunca es una garantía formal de orden cronológico, aunque coincida en este caso particular. Cómo corregirlo: siempre ordena explícitamente por la columna de tiempo antes de procesar eventos de forma incremental — sorted(events, key=lambda r: r[3]), como hace esta lección, nunca confíes en el orden del archivo fuente.
Usar INSERT también para el segundo y tercer evento de una sesión. Qué pasa: alguien, al portar el algoritmo de esta lección, olvida el if exists is None y usa INSERT ... ON CONFLICT o similar para todos los eventos, esperando que el conflicto de llave primaria "arregle" el problema silenciosamente. Por qué pasa: INSERT ... ON CONFLICT DO UPDATE es un patrón común en otros contextos, y puede parecer un atajo válido aquí. Cómo detectarlo: si tu conteo de filas finales no es 17, o si tu conteo de INSERT/UPDATE no suma exactamente 32, algo en la lógica de decisión está mal. Cómo corregirlo: el patrón de esta lección es explícito a propósito —SELECT para verificar existencia, luego INSERT o UPDATE según el resultado— porque hace visible, en el código mismo, la decisión central del accumulating snapshot: ¿esto es un proceso nuevo o uno que ya existe?
Pensar que el número de UPDATE (15) debería ser igual al número de sesiones con más de un evento. Qué pasa: alguien espera que 15 UPDATE corresponda a "15 sesiones que tuvieron más de un evento", cuando en realidad corresponde a "15 eventos que no fueron el primero de su sesión". Por qué pasa: es fácil confundir un conteo de eventos con un conteo de sesiones cuando ambos números aparecen juntos. Cómo detectarlo: si cuentas cuántas sesiones tienen add_to_cart_ts o purchase_ts no nulo (9 sesiones llegaron al carrito, algunas de esas también compraron), no vas a llegar a 15 — vas a llegar a un número menor, porque algunas sesiones generan dos UPDATE (una por add_to_cart, otra por purchase), no uno. Cómo corregirlo: 15 UPDATE es un conteo de eventos no iniciales, no de sesiones — cuenta cada add_to_cart y cada purchase por separado, incluso si ambos pertenecen a la misma sesión (como SESS-01, SESS-03, SESS-08, SESS-10, SESS-13, SESS-14, que generan dos UPDATE cada una).
Ejercicios
Ejercicio 1 — Cuenta cuántas sesiones generaron exactamente dos UPDATE. Sin volver a correr el script completo, usa fact_sessions ya construida para contar cuántas sesiones tienen tanto add_to_cart_ts como purchase_ts llenos (esas son, exactamente, las que generaron dos UPDATE cada una en el procesamiento incremental).
Ver solución
print(con.sql("""
SELECT COUNT(*) AS sessions_with_two_updates
FROM fact_sessions
WHERE add_to_cart_ts IS NOT NULL AND purchase_ts IS NOT NULL
"""))
Salida esperada:
┌───────────────────────────┐
│ sessions_with_two_updates │
│ int64 │
├───────────────────────────┤
│ 6 │
└───────────────────────────┘
Seis sesiones (SESS-01, SESS-03, SESS-08, SESS-10, SESS-13, SESS-14) generaron dos UPDATE cada una: 6 x 2 = 12 UPDATE. Las tres sesiones restantes que sí llegaron al carrito pero no compraron (SESS-07, SESS-11, SESS-17, del ejercicio 2 de la lección 3) generaron un solo UPDATE cada una: 3 x 1 = 3. 12 + 3 = 15, el total exacto de UPDATE que reportó el script de esta lección.
Ejercicio 2 — Simula un evento fuera de orden y observa qué pasa. Agrega, al final de events_sorted (sin reordenar), un evento purchase falso para una sesión SESS-99 que nunca tuvo un page_view previo. Corre el bucle de procesamiento sobre ese evento único y describe qué pasa.
Ver solución
con.execute("DELETE FROM fact_sessions WHERE session_id = 'SESS-99'") # por si ya existiera
exists = con.execute("SELECT 1 FROM fact_sessions WHERE session_id = 'SESS-99'").fetchone()
print(f"SESS-99 existe antes de procesar el evento huerfano: {exists}")
if exists is None:
# el evento es 'purchase', pero el codigo de esta leccion solo sabe INSERTAR
# sesiones nuevas cuando el PRIMER evento es un page_view -- un purchase
# huerfano no encaja en ninguna rama del if/elif, y se pierde en silencio
print("Este evento no coincide con ninguna rama del algoritmo: se descarta sin error visible.")
Salida esperada:
SESS-99 existe antes de procesar el evento huerfano: None
Este evento no coincide con ninguna rama del algoritmo: se descarta sin error visible.
El algoritmo de esta lección asume, siempre, que el primer evento de cualquier sesión es un page_view — no tiene ninguna rama que maneje "una sesión que empieza directamente con un purchase". Esto confirma, con evidencia, la advertencia de la profundización: un evento fuera de orden (o un milestone que llega sin su predecesor) no produce un error visible — simplemente no encaja en la lógica y se pierde. Resolver esto de forma robusta pertenece al patrón de dimensiones de llegada tardía, fuera del alcance de esta lección.
Ejercicio 3 — Explica, en tus propias palabras, por qué esta lección no necesitó cambiar ni una columna del esquema de fact_sessions respecto a la lección 3. En 2-3 frases, explica qué es exactamente lo que cambió entre las lecciones 3 y 4, y qué se mantuvo idéntico.
Ver solución
Lo que cambió fue cómo se llenan las columnas de fact_sessions —una consulta agregada de una sola vez (lección 3) contra un bucle de INSERT/UPDATE evento por evento en orden cronológico (lección 4)—, no qué columnas tiene la tabla ni qué significa cada una. El esquema (session_id, store_id, session_date, view_ts, add_to_cart_ts, purchase_ts, is_converted) es idéntico en ambas lecciones, porque ambas describen el mismo hecho de negocio con el mismo grano: una fila por sesión completa. Esta separación —el esquema define qué describe el hecho, el mecanismo de carga define cómo se llena— es una de las ideas centrales de todo modelado dimensional: el esquema sobrevive a un cambio de tecnología de carga (de un batch job nocturno a un consumidor de eventos en tiempo real, por ejemplo) sin que ninguna consulta que lea fact_sessions tenga que cambiar.
Resumen y siguiente paso
Esta lección reconstruyó fact_sessions procesando los 32 eventos de Kiosko uno a la vez, en orden cronológico —INSERT cuando una sesión empieza (17 veces), UPDATE cuando avanza (15 veces)—, y verificó, comparando fila por fila, que el resultado es idéntico al que produjo el enfoque agregado de la lección 3: cero diferencias. Esto confirma que el mecanismo mínimo de INSERT+UPDATE de la lección 2 escala sin cambios a la población completa de sesiones, y que es el mecanismo más cercano a como un pipeline de producción real mantendría esta tabla — recibiendo eventos de forma continua, sin esperar a tener todos los datos disponibles de antemano.
Antes de avanzar deberías poder: explicar por qué el orden cronológico de los eventos sí importa en esta lección aunque no importó en la lección 3; describir qué pasaría si un evento llegara fuera de orden; y calcular cuántos INSERT y UPDATE va a generar un conjunto de eventos, sin correr el código.
Con el primer patrón —accumulating snapshot— completo y verificado de dos formas distintas, la lección 5 cambia de tema: introduce el cumulative table design de Zach Wilson, un patrón que no describe un proceso con fin (como una sesión que termina en compra o abandono), sino una serie continua de actividad diaria por tienda, donde cada fila nueva se construye a partir de la fila del día anterior.
Recursos
- Kimball Group — "Accumulating Snapshot Fact Table" — el patrón que esta lección verifica con un segundo mecanismo de construcción. kimballgroup.com/data-warehouse-business-intelligence-resources/kimball-techniques/dimensional-modeling-techniques/accumulating-snapshot-fact-table. En inglés.
- Kimball Group — "Late Arriving Dimension" — la técnica relacionada con el problema de eventos fuera de orden que esta lección nombra en su profundización, ya introducida en el módulo 5 para dimensiones. kimballgroup.com/data-warehouse-business-intelligence-resources/kimball-techniques/dimensional-modeling-techniques/late-arriving-dimension. En inglés.
- DuckDB — documentación oficial del cliente Python, la interfaz que ejecutó cada
INSERT/UPDATEincremental de esta lección. duckdb.org/docs/current/clients/python/overview. En inglés.