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.
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.
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étodo | Funcionamiento | Por qué no se usa aquí |
|---|---|---|
| ClickPipes (fuente Snowflake) | Conector nativo de ClickHouse Cloud, sin ETL y con interfaz administrada | Snowflake 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 S3 | COPY INTO @stage exporta Parquet o CSV a S3; el conector S3 de ClickPipes lo carga en ClickHouse | Requiere 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-client | Utiliza 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 → ClickHouse | Un stream CDC de Snowflake alimenta un tema de Kafka; el conector Kafka de ClickPipes lo ingiere | Es 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 directamente | No 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.
--resumepermite 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:
- Detén el productor Snowflake para congelar datos.
- Ejecuta
python scripts/02_migrate_trips.py --resume: solo se transfieren las filas del delta, en segundos y no en minutos. - 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ón | Implementación | Motivo |
|---|---|---|
Motor trips_raw | ReplacingMergeTree(_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_trips | ReplacingMergeTree(updated_at) | Los viajes pueden corregirse, por ejemplo mediante ajustes de tarifa; updated_at es la columna de versión. |
Motor agg_hourly_zone_trips | ReplacingMergeTree(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 dbt | Es 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 = 8443El host y el puerto se guardan en .clickhouse_state. Carga el archivo en cualquier
terminal para incorporar los datos de conexión:
source .clickhouse_stateVerificación:
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/?query=SELECT+1" \
--user "default:${CLICKHOUSE_PASSWORD}"
# Expected: 1Paso 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.sqlTambié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: truenyc_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 ClickHouseA 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 pointResultado 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: ReplacingMergeTreeLas 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.pySalida 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: 1Estado 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.
Hoja 5: diseño de modelos dbt
Configura materialización, motor y estrategia incremental de cada modelo, con comentarios inmediatos.
04 Reconstrucción de la canalización dbt
Reconstruye Medallion en ClickHouse con dbt-clickhouse —incrementales delete_insert, ReplacingMergeTree y vistas actualizables— y crea el diccionario de zonas.