Ejemplo resuelto: un plan completo
Un plan de migración cumplimentado para la carga de trabajo de taxis de NYC, con el que comparar el tuyo una vez escrito.
Esta es la solución completa de las cinco hojas de trabajo
(1,
2,
3,
4,
5) aplicada a la carga
de trabajo NYC Taxi. Úsala para:
- Comprobar tus respuestas después de completar cada sección
- Entender el razonamiento de las decisiones que implementa la Parte 3
- Compararlas con la tabla de correspondencia de decisiones de la Parte 3 si elegiste
opciones distintas
Esta es la clave de respuestas: no la rellenes como si fuera tu plan. Completa migration-plan.md.
| Métrica | Valor |
|---|
| Total de tablas | 7 (TRIPS_RAW, FACT_TRIPS, AGG_HOURLY_ZONE_TRIPS, DIM_TAXI_ZONES, DIM_PAYMENT_TYPE, DIM_VENDOR, DIM_DATE) |
| Total de vistas | 2 (STG_TRIPS, STG_TAXI_ZONES) |
| Streams | 1 (TRIPS_CDC_STREAM en TRIPS_RAW) |
| Tareas | 2 (CDC_CONSUME_TASK, HOURLY_AGG_TASK) |
| Total de filas en TRIPS_RAW | ~50 000 000 |
| Intervalo de fechas | Ventana móvil de cuatro años que termina al configurar el entorno |
| Columnas VARIANT | 1 (TRIPS_RAW.TRIP_METADATA) |
| Usos de QUALIFY detectados | 1 (consulta Q3) |
| Usos de MERGE INTO detectados | 2 (HOURLY_AGG_TASK, CDC_CONSUME_TASK) |
| Objeto | Tipo | Esquema | Filas | Grado de complejidad | Notas |
|---|
trips_raw | Tabla | raw | ~50 M | B | RMT con columna de versión _synced_at; el solapamiento entre la carga masiva y CDC exige deduplicación; stg_trips debe usar FINAL |
stg_trips | Vista dbt | staging | — | B | JSONExtract para TRIP_METADATA; exige probar las rutas JSON |
stg_taxi_zones | Vista dbt | staging | — | A | Paso directo; trivial |
int_trips_enriched | Efímero dbt | staging | — | A | CTE; las diferencias SQL se resuelven en los modelos superiores |
fact_trips | Incremental dbt | analytics | ~50 M | C | Motor RMT; delete_insert; reescritura de QUALIFY; FINAL obligatorio |
agg_hourly_zone_trips | Incremental dbt | analytics | ~140 000 | B | RMT; ventana móvil de recálculo de 2 h; prueba con cuidado el límite de partición |
dim_taxi_zones | Tabla dbt | analytics | 265 | A | Referencia estática; recarga completa; trivial |
dim_payment_type | Tabla dbt | analytics | 6 | A | Referencia estática; trivial |
dim_vendor | Tabla dbt | analytics | 3 | A | Referencia estática; trivial |
taxi_zones_dict | Diccionario | analytics | 265 | B | Sintaxis específica de ClickHouse; dictGet() en tiempo de consulta |
mv_hourly_revenue | MV actualizable | analytics | — | B | Sintaxis REFRESH EVERY; comprobar la sustitución atómica |
TRIPS_CDC_STREAM / CDC_CONSUME_TASK | Stream + Task de Snowflake | — | — | D | Sin equivalente en ClickHouse; se sustituye por el cambio directo del productor en la Parte 3 |
| Tabla | Motor | Columna de versión | Razonamiento |
|---|
trips_raw | ReplacingMergeTree(_synced_at) | _synced_at | El script de migración de Python (scripts/02_migrate_trips.py) puede reintentar un lote y volver a insertar el mismo trip_id. Tras el cambio, el productor en vivo también puede reintentar ante fallos transitorios. _synced_at DateTime DEFAULT now() se define automáticamente al hacer INSERT: un reintento posterior tiene una marca de tiempo mayor, por lo que RMT conserva la escritura más reciente. stg_trips consulta con FINAL para imponer la deduplicación antes de ejecutar cualquier modelo posterior. |
fact_trips | ReplacingMergeTree(updated_at) | updated_at | Los viajes pueden corregirse (ajustes de tarifa, cambios de estado). El mismo trip_id se vuelve a insertar con valores actualizados. updated_at aumenta monótonamente con cada corrección: durante la deduplicación de RMT prevalece el valor mayor. Consulta siempre con FINAL. |
agg_hourly_zone_trips | ReplacingMergeTree(updated_at) | updated_at | dbt vuelve a calcular las últimas 2 horas y las reinserta. Sin RMT, los agregados antiguos y nuevos se acumulan y se cuentan dos veces. Definir updated_at como now() en cada ejecución de dbt garantiza que prevalezcan los valores más recientes. |
dim_taxi_zones | MergeTree() | — | Recarga completa mediante dbt (intercambio atómico de tablas, reconstrucción completa). No pueden acumularse duplicados. No hace falta deduplicar. |
dim_payment_type | MergeTree() | — | Igual: recarga completa. |
dim_vendor | MergeTree() | — | Igual: recarga completa. |
mv_hourly_revenue | MergeTree() | — | La MV actualizable sustituye atómicamente todo su conjunto de resultados en cada REFRESH. No hay upserts. |
| Tabla | ORDER BY | Razonamiento |
|---|
trips_raw | (pickup_at, trip_id) | Los recorridos por intervalos temporales filtran primero por pickup_at. trip_id es la clave de deduplicación de RMT y debe formar parte de ORDER BY para que RMT identifique las filas duplicadas. pickup_at va primero porque predominan los recorridos analíticos por intervalo; trip_id va al final porque tiene cardinalidad alta y solo actúa como discriminador de unicidad. |
fact_trips | (toStartOfMonth(pickup_at), pickup_at, trip_id) | Las siete consultas analíticas filtran por pickup_at. El prefijo mensual agrupa en bloques adyacentes los datos de cada mes natural, lo que permite omitir bloques de forma aproximada en las agregaciones mensuales sin añadir PARTITION BY. trip_id queda al final para aportar unicidad a RMT sin perjudicar la omisión de bloques. |
agg_hourly_zone_trips | (hour_bucket, zone_id) | Q6 (y todas las consultas de agregación) filtra por hour_bucket y zone_id. hour_bucket tiene unos 35 000 valores distintos; zone_id tiene 265. hour_bucket va primero porque los recorridos por intervalo temporal son el patrón de acceso principal. zone_id va después para el filtrado secundario. |
dim_taxi_zones | (location_id) | 265 filas equivalen a un gránulo. ORDER BY no afecta al rendimiento. Usar location_id como clave de unión es convencional y mejora la legibilidad. |
| Columna | Tipo de Snowflake | Tipo de ClickHouse | Justificación de la decisión |
|---|
TRIP_METADATA | VARIANT | String | Conserva exactamente el JSON original. JSONExtract* permite acceder a rutas arbitrarias durante la consulta. Map(String,String) pierde las estructuras anidadas; Tuple exige un esquema fijo. String es la elección segura para JSON arbitrario. |
PICKUP_DATETIME / PICKUP_AT | TIMESTAMP_NTZ(9) | DateTime64(3, 'UTC') | La precisión de milisegundos basta para las marcas de tiempo de viajes. Los nanosegundos (9) son excesivos. 'UTC' hace explícita la zona horaria y evita sorpresas relacionadas con el horario de verano en agregaciones por intervalo temporal. |
PICKUP_LOCATION_ID | INTEGER | UInt16 | Valores 1–265. El máximo de UInt8 es 255 (insuficiente). El máximo de UInt16 es 65535 (correcto). 2 bytes frente a 4 para Int32: ahorra unos 95 MB sin comprimir por columna en 50 millones de filas. |
VENDOR_ID | INTEGER | UInt8 | Valores 1–3. El máximo de UInt8 es 255: correcto. Un byte por fila. |
DRIVER_RATING | FLOAT | Nullable(Float32) | Es NULL con frecuencia (no todos los viajes tienen valoraciones). Nullable conserva la semántica correcta de nulos. Float32 basta para el intervalo 1,0–5,0. Float64 desperdiciaría almacenamiento sin aportar precisión útil. |
UPDATED_AT | TIMESTAMP_NTZ(9) | DateTime64(3, 'UTC') | Columna de versión para ReplacingMergeTree. Debe usar DateTime64, no DateTime: dos correcciones en el mismo segundo serían no deterministas con precisión de segundos. Los milisegundos garantizan el orden de deduplicación correcto. |
| Expresión de Snowflake | Equivalente de ClickHouse |
|---|
DATE_TRUNC('hour', pickup_at) | toStartOfHour(pickup_at) |
DATEADD('day', -7, CURRENT_DATE) | today() - 7 |
DATEDIFF('minute', pickup_at, dropoff_at) | dateDiff('minute', pickup_at, dropoff_at) |
TRIP_METADATA:driver.rating::FLOAT | JSONExtractFloat(trip_metadata, 'driver', 'rating') |
TRIP_METADATA:surge_multiplier::FLOAT | JSONExtractFloat(trip_metadata, 'surge_multiplier') |
QUALIFY ROW_NUMBER() OVER (PARTITION BY pickup_location_id ORDER BY fare_amount DESC) <= 10 | SELECT ... FROM (SELECT ..., ROW_NUMBER() OVER (...) AS rn FROM ...) WHERE rn <= 10 |
MERGE INTO fact_trips ... WHEN MATCHED THEN UPDATE | dbt incremental con delete_insert: DELETE de las filas con claves coincidentes y después INSERT de todas las filas nuevas |
| Oleada | Objetos | Dependencias | Notas |
|---|
| Oleada 0 | trips_raw (esquema), dim_taxi_zones, dim_payment_type, dim_vendor | Ninguna | dbt crea tablas vacías. Las tablas de dimensiones se rellenan de inmediato con datos de referencia estáticos (sin depender de los viajes). Ejecuta: dbt run --select trips_raw dim_* |
| Oleada 1 | Carga masiva en Python (scripts/02_migrate_trips.py) | Oleada 0 (el esquema de trips_raw debe existir) | 50 millones de filas de TRIPS_RAW en Snowflake. Se puede reanudar con --resume. Comprueba el recuento con scripts/01_verify_migration.sh. Unos 40-50 min. |
| Oleada 2 | stg_trips, stg_taxi_zones, int_trips_enriched, fact_trips, agg_hourly_zone_trips | Oleada 1 completada (trips_raw poblada) + Oleada 0 (existen las tablas de dimensiones) | dbt run completo. stg_trips lee trips_raw; int_trips_enriched se une a las dimensiones; fact_trips y agg_hourly_zone_trips se construyen encima. |
| Oleada 3 | taxi_zones_dict, mv_live_trip_feed | Oleada 2 (dim_taxi_zones poblada para el diccionario; fact_trips poblada para la MV) | Diccionario creado mediante scripts/04_create_dictionary.sql. MV actualizable creada con un modelo dbt; activar después su intervalo de actualización es un paso manual con ALTER TABLE ... MODIFY REFRESH, no algo que dbt ejecute automáticamente. |
| Oleada 4 | Cambio del productor (scripts/03_cutover.sh) | Oleada 1 completada (carga masiva verificada) + Oleada 2 completada (capa analítica construida) | Detén el productor de Snowflake; inicia el de ClickHouse para escribir directamente en ClickHouse Cloud; ejecuta dbt run para llenar agg_hourly_zone_trips con datos en vivo. |
| Objeto | Riesgo | Método de verificación |
|---|
fact_trips | Las consultas sin FINAL cuentan de más mientras se retrasan las fusiones. El intervalo de partición de delete_insert debe ceñirse al prefijo de ORDER BY para no eliminar particiones ajenas. | SELECT COUNT(*) FINAL coincide con Snowflake ± el retraso de CDC. Ejecuta dbt test. Compara los resultados de Q3 entre sistemas. Busca trip_ids duplicados: SELECT trip_id, count() FROM fact_trips GROUP BY trip_id HAVING count() > 1 LIMIT 10. |
agg_hourly_zone_trips | La ventana móvil de recálculo de 2 h debe limitar correctamente el intervalo de eliminación. Si es demasiado amplio, se borran agregados antiguos; si es demasiado estrecho, persisten agregados obsoletos. | Comprueba manualmente tuplas concretas (hour_bucket, zone_id) frente a Snowflake. Verifica que el trip_count total de todas las zonas coincida con AGG_HOURLY_ZONE_TRIPS de Snowflake para el mismo período. |
| Cambio del productor | Una interrupción del script de migración deja una brecha en el recuento de filas; vuelve a ejecutarlo con --resume para cubrirla. Un reintento del productor tras el cambio puede volver a insertar viajes ya presentes en ClickHouse. | scripts/01_verify_migration.sh: comprueba la paridad de recuentos entre Snowflake y ClickHouse. ReplacingMergeTree(_synced_at) gestiona de forma idempotente las inserciones duplicadas. |
Movimiento de datos: script de migración de Python (scripts/02_migrate_trips.py)
¿Por qué un script de Python en lugar de una transferencia mediante almacenamiento de objetos o ClickPipes?
remoteSecure() sirve para transferir datos de ClickHouse a ClickHouse; no se aplica aquí.
- La transferencia mediante almacenamiento de objetos (Snowflake → S3 → función de tabla S3 de ClickHouse) funcionaría, pero añade complejidad: exige aprovisionar un bucket S3, roles IAM y COPY INTO de Snowflake, una sobrecarga innecesaria para un laboratorio.
- ClickPipes no admite Snowflake como origen. Sus orígenes compatibles son Kafka, S3, Kinesis, CDC de PostgreSQL y CDC de MySQL.
- El script de Python usa
snowflake-connector-python y clickhouse-connect, paquetes ya instalados para el laboratorio. Muestra el progreso en tiempo real, permite reanudar con --resume tras una interrupción y su código puede inspeccionarse por completo.
Estrategia incremental (dbt): delete_insert
¿Por qué delete_insert en vez de las estrategias append o merge?
append inserta filas nuevas sin tocar las existentes. En fact_trips, donde las filas pueden actualizarse, genera duplicados. Es incorrecto.
merge (si estuviera disponible) sería lo más parecido a MERGE INTO de Snowflake, pero la estrategia merge de dbt-clickhouse presenta limitaciones con ReplacingMergeTree y no es la opción recomendada.
delete_insert elimina las filas del intervalo de claves del lote entrante y después inserta todas las filas nuevas. Es idempotente (al repetirlo produce el mismo resultado), gestiona tanto inserciones como actualizaciones y funciona correctamente con ReplacingMergeTree. Es la recomendación habitual de la comunidad de dbt-clickhouse para patrones de upsert.
| Criterio | Umbral | Medido mediante |
|---|
| Paridad del recuento de filas | Coincidencia ≥ 99,9 % (se espera CH ≥ SF tras el cambio) | scripts/01_verify_migration.sh |
| Paridad de checksum | Coincidencia MD5 en una muestra de 10 000 filas | scripts/02_validate_parity.sql |
| Tasa de éxito de pruebas dbt | 100 % | dbt test en dbt/nyc_taxi_dbt_ch |
| Paridad de resultados | Las 7 consultas devuelven los mismos resultados (dentro de la tolerancia de punto flotante) | Comparación manual en la salida de scripts/run_benchmark.sh |
| Modelo | Materialización | Motivo |
|---|
stg_trips | view | Lee y limpia trips_raw; este modelo no recibe actualizaciones; coste de almacenamiento cero; siempre refleja el estado actual del origen |
stg_taxi_zones | view | Igual: limpieza por paso directo de una tabla de origen |
int_trips_enriched | ephemeral | Lógica de unión pura que solo usa fact_trips; insertarla como CTE evita una tabla física redundante; ningún modelo la consulta directamente |
fact_trips | incremental | Los viajes pueden corregirse a posteriori; en cada ejecución solo deben procesarse las filas nuevas y actualizadas |
agg_hourly_zone_trips | incremental | El recálculo móvil de 2 horas es un patrón incremental: procesa las filas recientes, no los 50 millones |
dim_taxi_zones | table | 265 zonas estáticas; reconstrucción completa en cada ejecución de dbt mediante intercambio atómico de tablas; sin actualizaciones parciales |
dim_payment_type | table | 6 tipos estáticos; mismo razonamiento que para dim_taxi_zones |
dim_vendor | table | 3 proveedores; mismo razonamiento |
| Modelo | ENGINE | Columna de versión | Motivo |
|---|
fact_trips | ReplacingMergeTree(updated_at) | updated_at | Los viajes pueden corregirse; definir updated_at como now() en cada inserción hace que prevalezca la versión más reciente durante la deduplicación de fondo de RMT; delete_insert es la vía principal de corrección y RMT, la red de seguridad |
agg_hourly_zone_trips | ReplacingMergeTree(updated_at) | updated_at | El recálculo móvil reinserta agregados para los mismos pares (hour_bucket, zone_id); RMT garantiza que los agregados obsoletos se eliminen en la fusión de fondo |
dim_taxi_zones | MergeTree() | — | La recarga completa mediante dbt implica un intercambio atómico de tablas (reconstrucción completa) en cada ejecución; no pueden acumularse duplicados; no hace falta deduplicar |
dim_payment_type | MergeTree() | — | Igual que dim_taxi_zones |
dim_vendor | MergeTree() | — | Igual que dim_taxi_zones |
| Modelo | unique_key | incremental_strategy | Filtro incremental | ¿Por qué este filtro? |
|---|
fact_trips | trip_id | delete_insert | WHERE updated_at > (SELECT max(updated_at) FROM {{ this }}) | La marca de agua superior en updated_at captura tanto viajes nuevos como corregidos (los ajustes de tarifa vuelven a insertar el mismo trip_id con el mismo pickup_at, pero con un updated_at más reciente); una marca sobre pickup_at omitiría silenciosamente las correcciones |
agg_hourly_zone_trips | [hour_bucket, zone_id] | delete_insert | WHERE pickup_at >= now() - INTERVAL 2 HOUR | La ventana móvil de 2 horas obliga a reagregar las horas limítrofes para corregir siempre los recuentos de horas parciales; una marca superior en max(pickup_at) dejaría la hora limítrofe permanentemente infracontada |
| Modelo | ¿FINAL en la cláusula FROM? | Motivo |
|---|
stg_trips | Sí: FROM trips_raw FINAL | trips_raw es ReplacingMergeTree y puede contener filas trip_id duplicadas por reintentos del script de migración o del productor tras el cambio. stg_trips es el único punto de aplicación: deduplica aquí para que todos los modelos posteriores (int_trips_enriched, fact_trips, agg_hourly_zone_trips) reciban datos limpios |
int_trips_enriched | No | Lee de stg_trips (una vista), no de una tabla RMT; FINAL no es pertinente para vistas |
fact_trips | No (en el cuerpo del modelo) | delete_insert mantiene limpio fact_trips después de cada ejecución completada; añadir FINAL dentro del modelo lo aplicaría innecesariamente a la subconsulta is_incremental() que lee max(updated_at) de {{ this }}. Los paneles y las pruebas dbt usan FINAL externamente al consultar fact_trips de forma directa |
Este es el ejemplo completo. Tu migration-plan.md debe coincidir con las decisiones clave que aparecen aquí o documentar explícitamente por qué elegiste otras.