Snowflake MigrationClickHouse Workshops

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.


Lista de comprobación

  • Selección de motores: completada
  • Diseño de claves de ordenación: completado
  • Traducción del esquema: completada
  • Plan de oleadas de migración: completado
  • Diseño de modelos dbt: completado

Sección 1: resumen del perfil

MétricaValor
Total de tablas7 (TRIPS_RAW, FACT_TRIPS, AGG_HOURLY_ZONE_TRIPS, DIM_TAXI_ZONES, DIM_PAYMENT_TYPE, DIM_VENDOR, DIM_DATE)
Total de vistas2 (STG_TRIPS, STG_TAXI_ZONES)
Streams1 (TRIPS_CDC_STREAM en TRIPS_RAW)
Tareas2 (CDC_CONSUME_TASK, HOURLY_AGG_TASK)
Total de filas en TRIPS_RAW~50 000 000
Intervalo de fechasVentana móvil de cuatro años que termina al configurar el entorno
Columnas VARIANT1 (TRIPS_RAW.TRIP_METADATA)
Usos de QUALIFY detectados1 (consulta Q3)
Usos de MERGE INTO detectados2 (HOURLY_AGG_TASK, CDC_CONSUME_TASK)

Sección 2: inventario de objetos

ObjetoTipoEsquemaFilasGrado de complejidadNotas
trips_rawTablaraw~50 MBRMT con columna de versión _synced_at; el solapamiento entre la carga masiva y CDC exige deduplicación; stg_trips debe usar FINAL
stg_tripsVista dbtstaging—BJSONExtract para TRIP_METADATA; exige probar las rutas JSON
stg_taxi_zonesVista dbtstaging—APaso directo; trivial
int_trips_enrichedEfímero dbtstaging—ACTE; las diferencias SQL se resuelven en los modelos superiores
fact_tripsIncremental dbtanalytics~50 MCMotor RMT; delete_insert; reescritura de QUALIFY; FINAL obligatorio
agg_hourly_zone_tripsIncremental dbtanalytics~140 000BRMT; ventana móvil de recálculo de 2 h; prueba con cuidado el límite de partición
dim_taxi_zonesTabla dbtanalytics265AReferencia estática; recarga completa; trivial
dim_payment_typeTabla dbtanalytics6AReferencia estática; trivial
dim_vendorTabla dbtanalytics3AReferencia estática; trivial
taxi_zones_dictDiccionarioanalytics265BSintaxis específica de ClickHouse; dictGet() en tiempo de consulta
mv_hourly_revenueMV actualizableanalytics—BSintaxis REFRESH EVERY; comprobar la sustitución atómica
TRIPS_CDC_STREAM / CDC_CONSUME_TASKStream + Task de Snowflake——DSin equivalente en ClickHouse; se sustituye por el cambio directo del productor en la Parte 3

Sección 3: decisiones de selección de motores

TablaMotorColumna de versiónRazonamiento
trips_rawReplacingMergeTree(_synced_at)_synced_atEl 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_tripsReplacingMergeTree(updated_at)updated_atLos 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_tripsReplacingMergeTree(updated_at)updated_atdbt 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_zonesMergeTree()—Recarga completa mediante dbt (intercambio atómico de tablas, reconstrucción completa). No pueden acumularse duplicados. No hace falta deduplicar.
dim_payment_typeMergeTree()—Igual: recarga completa.
dim_vendorMergeTree()—Igual: recarga completa.
mv_hourly_revenueMergeTree()—La MV actualizable sustituye atómicamente todo su conjunto de resultados en cada REFRESH. No hay upserts.

Sección 4: diseño de claves de ordenación

TablaORDER BYRazonamiento
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.

Sección 5: notas sobre la traducción del esquema

ColumnaTipo de SnowflakeTipo de ClickHouseJustificación de la decisión
TRIP_METADATAVARIANTStringConserva 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_ATTIMESTAMP_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_IDINTEGERUInt16Valores 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_IDINTEGERUInt8Valores 1–3. El máximo de UInt8 es 255: correcto. Un byte por fila.
DRIVER_RATINGFLOATNullable(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_ATTIMESTAMP_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.

Traducciones de funciones necesarias

Expresión de SnowflakeEquivalente 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::FLOATJSONExtractFloat(trip_metadata, 'driver', 'rating')
TRIP_METADATA:surge_multiplier::FLOATJSONExtractFloat(trip_metadata, 'surge_multiplier')
QUALIFY ROW_NUMBER() OVER (PARTITION BY pickup_location_id ORDER BY fare_amount DESC) <= 10SELECT ... FROM (SELECT ..., ROW_NUMBER() OVER (...) AS rn FROM ...) WHERE rn <= 10
MERGE INTO fact_trips ... WHEN MATCHED THEN UPDATEdbt incremental con delete_insert: DELETE de las filas con claves coincidentes y después INSERT de todas las filas nuevas

Sección 6: oleadas de migración

OleadaObjetosDependenciasNotas
Oleada 0trips_raw (esquema), dim_taxi_zones, dim_payment_type, dim_vendorNingunadbt 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 1Carga 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 2stg_trips, stg_taxi_zones, int_trips_enriched, fact_trips, agg_hourly_zone_tripsOleada 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 3taxi_zones_dict, mv_live_trip_feedOleada 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 4Cambio 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.

Registro de riesgos (objetos de grado C/D)

ObjetoRiesgoMétodo de verificación
fact_tripsLas 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_tripsLa 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 productorUna 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.

Sección 7: diferencias conocidas entre dialectos

  • QUALIFY: afecta a Q3 (queries/q03_top_trips_qualify.sql)
  • Ruta con dos puntos de VARIANT: afecta a Q4 y Q5 (acceso al JSON de TRIP_METADATA)
  • LATERAL FLATTEN: no se usa en esta carga; se accede a VARIANT mediante una ruta con dos puntos, no con FLATTEN
  • MERGE INTO: afecta a los modelos incrementales dbt (fact_trips, agg_hourly_zone_trips)
  • Streams de Snowflake → cambio del productor (las escrituras en vivo van directamente a ClickHouse tras el cambio)
  • Diferencias en funciones de fecha: afectan a Q1 (DATE_TRUNC), Q3 (DATEADD) y Q4 (DATEDIFF)

Sección 8: estrategia de migración

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.

Sección 9: criterios para el cambio de sistema

CriterioUmbralMedido mediante
Paridad del recuento de filasCoincidencia ≥ 99,9 % (se espera CH ≥ SF tras el cambio)scripts/01_verify_migration.sh
Paridad de checksumCoincidencia MD5 en una muestra de 10 000 filasscripts/02_validate_parity.sql
Tasa de éxito de pruebas dbt100 %dbt test en dbt/nyc_taxi_dbt_ch
Paridad de resultadosLas 7 consultas devuelven los mismos resultados (dentro de la tolerancia de punto flotante)Comparación manual en la salida de scripts/run_benchmark.sh


Sección 10: diseño de modelos dbt

Selección de materialización

ModeloMaterializaciónMotivo
stg_tripsviewLee y limpia trips_raw; este modelo no recibe actualizaciones; coste de almacenamiento cero; siempre refleja el estado actual del origen
stg_taxi_zonesviewIgual: limpieza por paso directo de una tabla de origen
int_trips_enrichedephemeralLó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_tripsincrementalLos viajes pueden corregirse a posteriori; en cada ejecución solo deben procesarse las filas nuevas y actualizadas
agg_hourly_zone_tripsincrementalEl recálculo móvil de 2 horas es un patrón incremental: procesa las filas recientes, no los 50 millones
dim_taxi_zonestable265 zonas estáticas; reconstrucción completa en cada ejecución de dbt mediante intercambio atómico de tablas; sin actualizaciones parciales
dim_payment_typetable6 tipos estáticos; mismo razonamiento que para dim_taxi_zones
dim_vendortable3 proveedores; mismo razonamiento

Configuración de motores

ModeloENGINEColumna de versiónMotivo
fact_tripsReplacingMergeTree(updated_at)updated_atLos 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_tripsReplacingMergeTree(updated_at)updated_atEl 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_zonesMergeTree()—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_typeMergeTree()—Igual que dim_taxi_zones
dim_vendorMergeTree()—Igual que dim_taxi_zones

Estrategia incremental

Modelounique_keyincremental_strategyFiltro incremental¿Por qué este filtro?
fact_tripstrip_iddelete_insertWHERE 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_insertWHERE pickup_at >= now() - INTERVAL 2 HOURLa 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

Ubicación de FINAL

Modelo¿FINAL en la cláusula FROM?Motivo
stg_tripsSí: FROM trips_raw FINALtrips_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_enrichedNoLee de stg_trips (una vista), no de una tabla RMT; FINAL no es pertinente para vistas
fact_tripsNo (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.

En esta página

ES