Snowflake MigrationClickHouse Workshops

03 Aprovisionamiento y migración

Provisiona ClickHouse Cloud con Terraform, crea las tablas desde tu plan y mueve 50 millones de filas con un script Python reanudable.

Punto de partida

Módulo 02 completado: migration-plan.md está rellenado y todas las casillas de su lista de finalización están marcadas, mientras que el productor de Snowflake sigue en ejecución. El setup.sh de este módulo comprueba ese archivo y avisa si falta o está incompleto, pero nunca bloquea el avance. Nada de lo que hay aquí te impide continuar sin él; lo único que se resentirá será tu comprensión de los dos módulos siguientes. Reserva unos 60 minutos en total, de los cuales aproximadamente 40–50 corresponden a una transferencia de datos desatendida que puedes dejar ejecutándose en segundo plano. Aquí también comienza el consumo del crédito de prueba de ClickHouse Cloud: aprovisionar el servicio y completar este módulo cuesta alrededor de 1–2 dólares de crédito; el laboratorio entero cuesta unos 2–4 dólares.

Por qué

Este módulo convierte el plan en realidad. Todas las decisiones que anotaste en migration-plan.md durante el módulo 02 —qué motor MergeTree corresponde a cada tabla, qué clave ORDER BY se deriva de la carga real de consultas y cómo se traducen las construcciones exclusivas de Snowflake— se escriben directamente aquí en el DDL de las tablas; no se vuelven a deducir desde cero. ClickHouse no ofrece un índice que puedas añadir a posteriori: si una clave ORDER BY resulta incorrecta cuando ya hay 50 millones de filas en la tabla, la solución es una recarga completa, no un ALTER rápido.

Por eso también importa el control previo, aunque sea flexible y no pueda detenerte. Si ejecutas este módulo sin un plan terminado, el proceso seguirá funcionando de forma mecánica: dbt run creará fact_trips como ReplacingMergeTree y el script de migración moverá los 50 millones de filas. Sin embargo, no sabrás por qué se eligió ese motor y no un MergeTree sencillo, por qué la clave de ordenación tiene esa forma ni cómo justificar las aceleraciones de aproximadamente 6–9 veces que mostrará después el módulo 04. La tabla de correspondencia de decisiones que aparece más abajo relaciona cada elección aplicada aquí con la pregunta de la hoja de trabajo que responde, para que puedas contrastarla con tu propio plan antes de aprovisionar nada.

Conceptos internos

Arquitectura de destino. Snowflake continúa recibiendo viajes nuevos mediante el productor mientras un script Python de ejecución única rellena ClickHouse con los 50 millones de filas existentes. Los dos sistemas funcionan en paralelo durante toda la migración; esto todavía no es la transición.

Flujo de migración: el productor sigue escribiendo en TRIPS_RAW mientras un script mueve 50 millones de filas en lotes de 100.000 a ClickHouse Cloud

En ClickHouse, trips_raw es la tabla de aterrizaje en la que escribe el script. A partir de ella, dbt construye las vistas de staging y el resto de la capa de analytics: este módulo crea ese esquema, pero todavía no lo rellena más allá de trips_raw.

Destino ClickHouse: trips_raw en ReplacingMergeTree alimenta vistas staging, hechos, dimensiones, agregado horario y diccionario de zonas

Leyenda:

  • Verde: fuentes de datos, antes y después de la transición.
  • Azul: tablas Snowflake.
  • Naranja: modelos y canalización dbt.
  • Rojo: tablas y vistas materializadas ClickHouse.
  • Cian: paneles Apache Superset.
  • Flechas discontinuas: flujos posteriores a la transición.

Por qué se usa un script Python y no un conector nativo. Existen varios métodos para trasladar datos de Snowflake a ClickHouse. Este laboratorio utiliza un script Python por lotes; la tabla siguiente explica por qué frente a las alternativas:

MétodoFuncionamientoPor qué no se usa aquí
ClickPipes (fuente Snowflake)Conector nativo de ClickHouse Cloud, sin ETL y con interfaz administradaSnowflake no es una fuente compatible con ClickPipes. ClickPipes admite Kafka, S3, Kinesis, CDC de PostgreSQL, CDC de MySQL y almacenamiento de objetos.
Exportación S3 → ClickPipes S3COPY INTO @stage exporta Parquet o CSV a S3; el conector S3 de ClickPipes lo carga en ClickHouseRequiere un bucket S3, un rol IAM, un stage de Snowflake y una cuenta de AWS. Añade unas tres etapas de preparación antes de mover un solo dato. Es viable en producción, pero supone demasiada infraestructura para un laboratorio.
Exportación S3 → clickhouse-clientUtiliza la misma exportación a S3, pero la carga con INSERT INTO ... SELECT FROM s3(...)Tiene los mismos requisitos previos de S3. Además, obliga al participante a gestionar manualmente la fragmentación de archivos y la reanudación.
Snowflake → Kafka → ClickHouseUn stream CDC de Snowflake alimenta un tema de Kafka; el conector Kafka de ClickPipes lo ingiereEs una canalización de streaming completa, adecuada para requisitos de latencia inferiores a un minuto en producción. Un clúster de Kafka resulta excesivo para un entorno de laboratorio.
Script Python (este laboratorio)snowflake-connector-python lee el cursor en lotes de 100 000 filas; clickhouse-connect inserta directamenteNo requiere infraestructura adicional aparte de los paquetes que ya necesita el laboratorio. Se puede reanudar con --resume, usando max(pickup_at) como marca de agua. Muestra el progreso en tiempo real. Tarda unos 40–50 minutos para 50 millones de filas a unas 20 000 filas/s, aceptable para un ejercicio de migración puntual.

Por qué el script Python es la opción adecuada para este laboratorio:

  • No requiere una cuenta de AWS. Los enfoques basados en S3 obligan a crear un bucket, políticas IAM y un stage externo de Snowflake: tres pasos que no tienen nada que ver con ClickHouse.
  • Es autocontenido. Los dos paquetes (snowflake-connector-python, clickhouse-connect) se instalan en el mismo entorno virtual que dbt. No hacen falta servicios ni credenciales nuevos.
  • Se puede reanudar. --resume permite interrumpir y reiniciar el script con seguridad. ReplacingMergeTree(_synced_at) garantiza que las inserciones duplicadas de un reintento se dedupliquen automáticamente.
  • Es transparente. Los participantes pueden leer el script, entender la correspondencia entre columnas y adaptarlo a su propio esquema, algo más educativo que avanzar por un asistente gráfico.

Tratamiento del desfase de migración. El productor de Snowflake sigue activo durante los aproximadamente 40–50 minutos que tarda el script. Los viajes escritos en Snowflake durante esa ventana todavía no están en ClickHouse. El laboratorio cierra el desfase con un enfoque de dos pasadas durante la transición, que el módulo 05 recorre de forma explícita:

  1. Detén el productor Snowflake para congelar datos.
  2. Ejecuta python scripts/02_migrate_trips.py --resume: solo se transfieren las filas del delta, en segundos y no en minutos.
  3. Inicia el productor ClickHouse.

La misma deduplicación ReplacingMergeTree(_synced_at) que protege los reintentos de la migración también resuelve este caso: si alguna fila se solapa entre la ejecución de este módulo y la pasada posterior con --resume, gana el valor _synced_at más reciente.

Cuándo elegir S3 en producción. Si el conjunto de datos supera los 500 millones de filas o si el coste de consultar todo el warehouse de Snowflake resulta significativo, es preferible exportar mediante S3: Snowflake genera Parquet comprimido en paralelo, mucho más rápido que un único cursor, y ClickHouse también puede cargar desde S3 en paralelo. El enfoque con Python funciona bien a la escala del laboratorio.

Correspondencia de decisiones. La tabla siguiente contiene la misma lista de decisiones de las hojas 1 (selección del motor), 2 (claves de ordenación) y 3 (traducción del esquema) de migration-plan.md, contrastada con lo que realmente construye el laboratorio. Compárala con tu propio plan antes de aprovisionar nada:

DecisiónImplementaciónMotivo
Motor trips_rawReplacingMergeTree(_synced_at)El script Python de migración usa inserciones por lotes que pueden repetirse tras una interrupción. _synced_at DateTime DEFAULT now() se define en cada INSERT, por lo que una fila reintentada llega después y tiene un valor _synced_at superior; la fila posterior gana en la deduplicación RMT y los reintentos son idempotentes. Los reintentos del productor posteriores a la transición son seguros por el mismo motivo. stg_trips consulta con FINAL para garantizar una fila por viaje.
Motor fact_tripsReplacingMergeTree(updated_at)Los viajes pueden corregirse, por ejemplo mediante ajustes de tarifa; updated_at es la columna de versión.
Motor agg_hourly_zone_tripsReplacingMergeTree(updated_at)El recálculo de una ventana móvil sigue un patrón de upsert.
Tablas dim_*MergeTree()Se recargan por completo en cada ejecución de dbt; no hay upserts.
ORDER BY de fact_trips(toStartOfMonth(pickup_at), pickup_at, trip_id)Todas las consultas Q1–Q7 filtran por pickup_at; trip_id garantiza la unicidad en el nivel más específico.
ORDER BY de agg_hourly_zone_trips(hour_bucket, zone_id)Ambas columnas aparecen en todas las consultas de agregación.
VARIANT →String + JSONExtract*Conserva el JSON sin procesar; la extracción se realiza al consultar.
QUALIFY →Subconsulta que envuelve ROW_NUMBER()ClickHouse dispone de una cláusula QUALIFY nativa desde la versión 24.5, pero se enseña la forma con subconsulta porque es portable a versiones de ClickHouse y motores SQL anteriores o sin QUALIFY.
MERGE INTO →Incremental delete_insert en dbtEs la estrategia de upsert idiomática de dbt-clickhouse y evita reescribir toda la tabla.

Paso 1 — Provisiona el clúster ClickHouse

setup.sh realiza una sola tarea: ejecuta terraform apply y escribe los datos de conexión en .clickhouse_state. También vuelve a comprobar que exista migration-plan.md antes de aprovisionar nada —consulta la sección Por qué—, pero esa comprobación solo avisa y nunca bloquea.

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"

# Configure credentials
cp .env.example .env
vim .env
# Fill in: CLICKHOUSE_ORG_ID, CLICKHOUSE_TOKEN_KEY, CLICKHOUSE_TOKEN_SECRET, CLICKHOUSE_PASSWORD

# Provision
source .env && ./setup.sh

.env está excluido de git; no lo añadas nunca al repositorio.

Salida esperada: Terraform crea 2 recursos, el servicio y la lista de acceso IP, en aproximadamente 2–3 minutos:

Apply complete! Resources: 2 added, 0 changed, 0 destroyed.

Outputs:
clickhouse_host = "abc123xyz.us-east-1.aws.clickhouse.cloud"
clickhouse_port = 8443

El host y el puerto se guardan en .clickhouse_state. Carga el archivo en cualquier terminal para incorporar los datos de conexión:

source .clickhouse_state

Verificación:

curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+1" \
  --user "default:${CLICKHOUSE_PASSWORD}"
# Expected: 1

Paso 2 — Crea las tablas vacías

Crea primero trips_raw manualmente con el motor correcto. El script de migración cargará datos en esta tabla durante el Paso 3; debe existir de antemano con ReplacingMergeTree para que la columna de versión esté configurada antes de que llegue la primera fila.

-- Run in the ClickHouse SQL console (cloud.clickhouse.com -> SQL console)
CREATE TABLE IF NOT EXISTS default.trips_raw (
    trip_id               String,
    vendor_id             UInt8,
    pickup_at             DateTime64(3, 'UTC'),
    dropoff_at            DateTime64(3, 'UTC'),
    passenger_count       UInt8,
    trip_distance_miles   Float32,
    pickup_location_id    UInt16,
    dropoff_location_id   UInt16,
    payment_type_id       UInt8,
    rate_code_id          UInt8,
    store_fwd_flag        String,
    fare_amount_usd       Float32,
    extra_amount_usd      Float32,
    mta_tax_usd           Float32,
    tip_amount_usd        Float32,
    tolls_amount_usd      Float32,
    total_amount_usd      Float32,
    ingested_at           DateTime64(3, 'UTC'),
    trip_metadata         String,
    _synced_at            DateTime DEFAULT now()
)
ENGINE = ReplacingMergeTree(_synced_at)
ORDER BY (pickup_at, trip_id);

_synced_at se define automáticamente en cada INSERT. Si el script de migración se interrumpe y vuelve a ejecutarse con --resume, durante un breve periodo pueden coexistir filas duplicadas para el mismo trip_id. RMT conserva la fila posterior, la que tiene el valor _synced_at más alto. stg_trips consulta trips_raw FINAL para forzar la deduplicación antes de que ningún modelo posterior vea los datos.

A continuación, carga los datos de referencia de las zonas. Son datos estáticos, las 265 zonas de la NYC TLC, que el modelo dbt stg_taxi_zones lee como fuente.

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state

clickhouse-client --host "${CLICKHOUSE_HOST}" --port 9440 --secure \
  --user default --password "${CLICKHOUSE_PASSWORD}" \
  --multiquery < scripts/00_seed_zones.sql

También puedes pegar directamente el contenido de scripts/00_seed_zones.sql en la consola SQL de ClickHouse.

Configura el perfil dbt. El dbt_project.yml de este proyecto declara profile: 'nyc_taxi_ch'. Si no existe un perfil con ese nombre en ~/.dbt/profiles.yml, dbt run falla inmediatamente con Could not find profile named 'nyc_taxi_ch'. La plantilla se encuentra en workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch/profiles.yml.example.

El módulo 01 ya escribió en ~/.dbt/profiles.yml un perfil nyc_taxi: para Snowflake, y su bucle de actualización del Paso 4 continúa consultándolo mientras el productor de Snowflake siga activo. No sustituyas ese archivo por la plantilla de ClickHouse: sobrescribirlo con profiles.yml.example eliminaría el perfil nyc_taxi: y rompería el bucle de actualización del módulo 01. Abre la plantilla y combina su bloque nyc_taxi_ch: con el ~/.dbt/profiles.yml existente como segundo perfil de nivel superior, junto a nyc_taxi::

nyc_taxi:        # from module 01 — leave this one alone
  target: dev
  outputs:
    dev:
      type: snowflake
      # ...

nyc_taxi_ch:      # add this block
  target: dev
  outputs:
    dev:
      type: clickhouse
      schema: nyc_taxi_ch
      host: "{{ env_var('CLICKHOUSE_HOST') }}"
      port: 8443
      user: "{{ env_var('CLICKHOUSE_USER', 'default') }}"
      password: "{{ env_var('CLICKHOUSE_PASSWORD') }}"
      secure: true

nyc_taxi_ch: lee CLICKHOUSE_HOST, CLICKHOUSE_USER y CLICKHOUSE_PASSWORD del entorno mediante env_var(). Por eso debes cargar .env y .clickhouse_state antes de cualquier orden dbt de este módulo; el dbt run de más abajo ya lo hace. Igual que el perfil de Snowflake, ~/.dbt/profiles.yml contiene credenciales y está excluido de git: no lo añadas nunca al repositorio. Esta combinación agrega un segundo juego de credenciales a un archivo que ya contenía el primero.

Verificación:

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.clickhouse_state"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.env"

dbt debug
# Expected: "All checks passed!" — confirms dbt found the nyc_taxi_ch profile and
# connected to ClickHouse

A continuación, ejecuta dbt run para crear las tablas de analytics y las vistas de staging:

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.clickhouse_state"
source "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/.env"

dbt deps   # install packages (first run only)
dbt run    # creates analytics tables and staging views; all empty at this point

Resultado esperado: unos 8 modelos creados en menos de 2 minutos, con todas las tablas vacías.

Verificación:

# Check analytics tables were created
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SHOW+TABLES+IN+analytics" \
  --user "default:${CLICKHOUSE_PASSWORD}"
# Expected: agg_hourly_zone_trips, dim_date, dim_payment_type, dim_vendor, dim_taxi_zones, fact_trips

# Check staging views were created
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SHOW+TABLES+IN+staging" \
  --user "default:${CLICKHOUSE_PASSWORD}"
# Expected: stg_trips, stg_taxi_zones

# Check trips_raw exists with the correct engine
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+engine+FROM+system.tables+WHERE+database%3D%27default%27+AND+name%3D%27trips_raw%27" \
  --user "default:${CLICKHOUSE_PASSWORD}"
# Expected: ReplacingMergeTree

Las seis tablas de analytics y las dos vistas de staging ya existen, pero todas siguen vacías: dbt run solo ha creado sus esquemas. La única tabla que contendrá datos después de este módulo es trips_raw, aunque todavía tampoco tiene ninguno; eso ocurre a continuación.

Paso 3 — Migra los datos

Carga todas las filas de NYC_TAXI_DB.RAW.TRIPS_RAW en Snowflake dentro de default.trips_raw en ClickHouse mediante el script Python de migración por lotes.

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state
source .venv/bin/activate
python scripts/02_migrate_trips.py

Salida esperada (aproximadamente 40–50 minutos para 50 millones de filas):

━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━
  NYC Taxi Migration: Snowflake -> ClickHouse
━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━

  Rows to migrate: 50,000,000
  Batch size:      100,000

  Rows inserted    Elapsed      ETA                    Rate
  -------------------- ------------ ---------------------- ---------------
  100,000              0m 07s       56m 14s remaining      13,945 rows/s
  200,000              0m 14s       55m 28s remaining      14,021 rows/s
  ...

El productor de Snowflake sigue escribiendo en TRIPS_RAW durante toda la ejecución del script. Por tanto, ClickHouse queda retrasado aproximadamente lo que dura la transferencia. Ese desfase es intencionado y se corrige en el módulo 05, no aquí.

Si el script se interrumpe, vuelve a ejecutarlo con --resume para continuar desde el último punto de control:

python scripts/02_migrate_trips.py --resume

--resume lee max(pickup_at) en ClickHouse y omite las filas ya cargadas. Por eso siempre puedes interrumpir y reiniciar el script con seguridad: nunca terminarás con una carga parcial imposible de recuperar.

Cómo verificar la finalización

cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse"
source .env && source .clickhouse_state

# Row count in trips_raw
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+count()+FROM+default.trips_raw" \
  --user "default:${CLICKHOUSE_PASSWORD}"
# Expected: approximately 50000000

# .clickhouse_state was written by setup.sh
ls -la .clickhouse_state

# Service is reachable
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+1" \
  --user "default:${CLICKHOUSE_PASSWORD}"
# Expected: 1

Estado final

El servicio de ClickHouse Cloud está activo y accesible, y setup.sh ha escrito en disco .clickhouse_state con CLICKHOUSE_HOST y CLICKHOUSE_PORT. Existen todas las tablas de destino y las vistas de staging: default.trips_raw, las dos vistas (stg_trips, stg_taxi_zones), las seis tablas analytics (fact_trips, agg_hourly_zone_trips, dim_taxi_zones, dim_payment_type, dim_vendor, dim_date) y la vista materializada actualizable analytics.mv_live_trip_feed. En total, el dbt run de este módulo creó siete objetos en el esquema analytics. Todas las tablas de analytics siguen vacías excepto mv_live_trip_feed, que ya contiene la única fila de instantánea que produjo dbt run al construir la vista. Por lo demás, solo default.trips_raw contiene datos: aproximadamente 50 millones de filas trasladadas por el script Python del Paso 3.

El productor de Snowflake sigue en ejecución. Nunca se ha detenido en este módulo y no debe detenerse aquí. Cada viaje que se escriba en el TRIPS_RAW de Snowflake después del último lote del script de migración es una fila que ClickHouse aún no posee. Por eso, ClickHouse queda retrasado respecto a Snowflake aproximadamente lo que dura la ventana de migración, unos 40–50 minutos, más el tiempo de preparación del módulo. El desfase es real y sigue creciendo mientras el productor permanezca activo. No lo cierres en este módulo. La transición del módulo 05 lo cierra de forma deliberada en un paso controlado de dos pasadas, que mide el tamaño del desfase antes de eliminarlo. Detener el productor o volver a ejecutar ahora el script de migración suprimiría precisamente aquello que el módulo 05 está diseñado para demostrar.

En esta página

¿Quieres seguir tu progreso?

Opcional. Enviaremos un enlace por correo para confirmar tu dirección; el progreso se registrará cuando lo abras.

Usa tu correo de trabajo, no uno personal.

Para seguir el progreso también debes aceptar los Términos del servicio actuales en la Configuración de privacidad.

ES