Snowflake MigrationClickHouse Workshops

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.

Punto de partida

Módulo 03 completado: el servicio de ClickHouse Cloud está activo y accesible, y setup.sh ha escrito .clickhouse_state en disco 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; el dbt run del módulo 03 ya construyó los siete objetos. Todas las tablas analytics siguen vacías excepto mv_live_trip_feed, que ya contiene la única fila de instantánea de esa construcción. Por lo demás, solo default.trips_raw tiene datos: aproximadamente 50 millones de filas. El productor de Snowflake sigue en ejecución, por lo que ClickHouse lleva un retraso similar a la duración de la ventana de migración. Reserva unos 30 minutos.

Por qué

El módulo 03 demostró que ClickHouse puede almacenar 50 millones de filas. No demostró que la canalización pueda ejecutarse en ClickHouse: las vistas de staging, la tabla de hechos incremental, las recargas de dimensiones y las pruebas que detectan un modelo defectuoso antes de que lo vea un participante. Eso es lo que reconstruye este módulo: los mismos modelos Medallion del módulo 01, expresados con dbt-clickhouse en lugar de dbt-snowflake y ejecutados sobre las tablas que creó el módulo 03.

La lógica de los modelos no cambia: stg_trips continúa convirtiendo tipos y extrayendo JSON, int_trips_enriched sigue uniendo las dimensiones y fact_trips sigue terminando con una fila por viaje. Lo que cambia es la capa de materialización subyacente: no hay MERGE INTO, ni Task de Snowflake, ni cluster_by. Este módulo demuestra que la migración no es un volcado de datos puntual: la canalización que el equipo de un participante ejecuta a diario, con la misma programación y protegida por la misma ejecución de dbt test, continúa funcionando cuando cambia el warehouse que la sustenta.

Conceptos internos

El conjunto de modelos difiere de la canalización de Snowflake en cuatro aspectos. La referencia de configuración completa se encuentra en dbt en ClickHouse; esta es la versión breve que necesitas antes de ejecutar dbt run en el Paso 1. Para consultar la canalización de origen que sustituyen estos modelos, consulta dbt en Snowflake.

1. delete_insert sustituye a MERGE. ClickHouse no tiene una instrucción MERGE INTO. Allí donde la canalización de Snowflake usaba incremental_strategy: merge para realizar upserts en fact_trips y agg_hourly_zone_trips, los modelos de ClickHouse emplean incremental_strategy: delete_insert: dbt elimina las filas cuya unique_key coincide con el lote entrante y después inserta el lote. En fact_trips, la unique_key es trip_id, y el filtro incremental usa updated_at como marca de agua en lugar de pickup_at. Una corrección de tarifa vuelve a insertar el mismo trip_id con el mismo pickup_at, pero con un updated_at más reciente; una marca de agua basada en pickup_at omitiría silenciosamente esa corrección.

2. ReplacingMergeTree es la red de seguridad bajo delete_insert, no su sustituto. Los dos modelos incrementales se declaran como ReplacingMergeTree(updated_at). Si una ejecución de delete_insert termina normalmente, la tabla ya contiene una fila por clave y el motor no tiene nada que limpiar. Si una ejecución se interrumpe a mitad de camino —después de eliminar y antes de insertar—, las fusiones en segundo plano terminan deduplicando las filas restantes y conservan la que tenga el valor updated_at más alto. Nunca confíes solo en ReplacingMergeTree para realizar el trabajo de deduplicación que corresponde a delete_insert: las fusiones en segundo plano son asíncronas y pueden retrasarse de minutos a horas en una tabla de este tamaño.

3. Las vistas materializadas actualizables sustituyen a la tarea programada. La canalización de Snowflake usaba una Task programada que ejecutaba un procedimiento almacenado para mantener actualizado un agregado móvil. En cambio, el proyecto dbt de ClickHouse declara mv_live_trip_feed con materialized = 'materialized_view' y engine = 'ReplacingMergeTree(refreshed_at)'. dbt run la construye como vista materializada actualizable, frente a las más de 30 líneas de DDL CREATE TASK que producían el mismo efecto en Snowflake. Este módulo no activa ningún intervalo de actualización; hacerlo requeriría una instrucción manual ALTER TABLE analytics.mv_live_trip_feed MODIFY REFRESH EVERY ... que el laboratorio no automatiza (el módulo 05 explica por qué).

4. mv_live_trip_feed no tiene ningún equivalente en Snowflake. No es la traducción de un modelo existente, sino una capacidad nueva que añade la migración. Una vista materializada estándar de ClickHouse se activa una vez por INSERT y solo ve las filas de ese lote; por tanto, no puede calcular correctamente un agregado histórico completo, como el total de viajes o la tarifa media. Una vista materializada REFRESHABLE vuelve a ejecutar toda su consulta —en este caso, SELECT ... FROM {{ ref('fact_trips') }}— según un calendario, de modo que cada actualización ve la tabla completa. El lado Snowflake de este taller nunca contó con esta opción.

Estos son los modelos que realmente construye dbt y el momento en que cada uno recibe datos:

ModeloCapaMaterializaciónPoblaciónNotas
stg_tripsstagingVistaPaso 1 (cada ejecución)Convierte tipos, JSONExtract* para trip_metadata
stg_taxi_zonesstagingVistaPaso 1 (cada ejecución)Lectura directa de la dimensión de zonas
int_trips_enrichedstagingEfímera— (incorporada como CTE)Todas las uniones de dimensiones; sin tabla física
fact_tripsanalyticsIncrementalPaso 1delete_insert por trip_id, marca de agua updated_at; ORDER BY (toStartOfMonth(pickup_at), pickup_at, trip_id)
agg_hourly_zone_tripsanalyticsIncrementalMódulo 05, después de la transiciónVentana móvil de 2 horas; el filtro incremental solo coincide con filas del productor en vivo; consulta el Paso 1
dim_taxi_zonesanalyticsTablaPaso 1Recarga completa en cada ejecución; fuente del diccionario de zonas del Paso 2
dim_payment_typeanalyticsTablaPaso 1Recarga completa en cada ejecución
dim_vendoranalyticsTablaPaso 1Recarga completa en cada ejecución
dim_dateanalyticsTablaPaso 1Calendario estático de 2009–2029; recarga completa en cada ejecución
mv_live_trip_feedanalyticsVista materializada (actualizable)dbt run del módulo 03; nunca se activa el intervalo de actualizaciónNo tiene equivalente en Snowflake; consulta el punto 4 anterior

Paso 1 — Puebla analytics

Activa el entorno virtual de dbt-clickhouse que creaste en el módulo 00 y ejecuta dbt run por segunda vez. El módulo 03 ya lo ejecutó una vez sobre tablas vacías para crear los esquemas; esta ejecución dispone de datos reales: trips_raw contiene ahora 50 millones de filas.

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

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

dbt lee la conexión de ClickHouse desde ~/.dbt/profiles.yml, a partir de dbt/nyc_taxi_dbt_ch/profiles.yml.example: es el mismo perfil que utilizó el módulo 03 para crear los esquemas vacíos.

Resultado esperado: aproximadamente 8–12 minutos, mientras los modelos incrementales procesan 50 millones de filas.

agg_hourly_zone_trips quedará vacía después de esta ejecución; es el resultado esperado, no un fallo. Su filtro incremental es WHERE pickup_at >= now() - INTERVAL 2 HOUR, que solo coincide con filas escritas por el productor en vivo. Todas las filas que acabas de migrar son históricas, por lo que ninguna cae dentro de una ventana de dos horas medida desde este momento. La tabla permanece vacía hasta que la transición del módulo 05 inicie el productor de ClickHouse; no pierdas tiempo intentando depurarla como si la canalización estuviera rota.

A continuación, ejecuta el conjunto de pruebas:

dbt test

Resultado esperado: todas las pruebas pasan.

Verificación:

SELECT formatReadableQuantity(count()) FROM analytics.fact_trips FINAL;
-- Expected: ~50 million

SELECT count() FROM analytics.agg_hourly_zone_trips;
-- Expected: 0 (normal — populated after cutover in module 05)

SELECT count() FROM analytics.dim_taxi_zones;
-- Expected: 265

Paso 2 — Crea el diccionario de zonas

analytics.dim_taxi_zones ya contiene las 265 zonas de la NYC TLC. Crea analytics.taxi_zones_dict, un diccionario en memoria respaldado por esa tabla, para que las consultas posteriores puedan obtener el distrito de una zona mediante dictGet() en lugar de un JOIN.

Qué aporta un diccionario frente a una unión. El diccionario se carga en memoria una vez y permanece listo; cada búsqueda posterior tiene un coste prácticamente nulo. Un JOIN con dim_taxi_zones vuelve a leer y emparejar la tabla de dimensiones en cada ejecución. Para una tabla de referencia pequeña y poco cambiante como esta —265 filas, recargadas por completo con cada dbt run— la elección es clara. El benchmark del módulo 05 consulta taxi_zones_dict directamente mediante dictGet, por lo que este paso es una dependencia obligatoria de ese módulo, no un complemento opcional.

Carga los datos de conexión y después aplica el DDL del diccionario mediante clickhouse-client o la API HTTP; elige la opción que tengas disponible:

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

# Via clickhouse-client
clickhouse-client \
  --host "${CLICKHOUSE_HOST}" \
  --port 9440 \
  --user default \
  --password "${CLICKHOUSE_PASSWORD}" \
  --secure \
  --multiquery \
  < scripts/04_create_dictionary.sql

# Or via HTTP API
curl "https://${CLICKHOUSE_HOST}:${CLICKHOUSE_PORT}/" \
  --user "default:${CLICKHOUSE_PASSWORD}" \
  --data-binary @scripts/04_create_dictionary.sql

Verificación:

-- Should return 'Manhattan' for zone 42
SELECT dictGet('analytics.taxi_zones_dict', 'borough', toUInt16(42));

-- Should show status = LOADED, element_count = 265
SELECT name, status, element_count
FROM system.dictionaries
WHERE name = 'taxi_zones_dict';

Cómo verificar la finalización

SELECT formatReadableQuantity(count()) FROM analytics.fact_trips FINAL;
-- Expected: ~50 million
SELECT count() FROM analytics.agg_hourly_zone_trips;
-- Expected: 0

Cero es correcto y no indica ningún fallo. El filtro incremental de agg_hourly_zone_trips (WHERE pickup_at >= now() - INTERVAL 2 HOUR) solo coincide con filas escritas por el productor en vivo. Todas las filas que hay ahora en ClickHouse son datos históricos trasladados por el script de migración del módulo 03; ninguna tiene menos de dos horas respecto a now(). La tabla solo se llena cuando la transición del módulo 05 inicia el productor de ClickHouse. Hasta entonces, cualquier gráfico de panel basado en esta tabla no mostrará datos, y ese comportamiento es esperado, no algo que debas depurar.

SELECT count() FROM analytics.dim_taxi_zones;
-- Expected: 265
cd "$(git rev-parse --show-toplevel)/workshop_public/snowflake_migration_lab/03-migrate-to-clickhouse/dbt/nyc_taxi_dbt_ch"
dbt test

Esperado: todos pasan.

-- Should return a borough name, e.g. 'Manhattan'
SELECT dictGet('analytics.taxi_zones_dict', 'borough', toUInt16(42));

Estado final

La capa de analytics está poblada y probada: fact_trips contiene aproximadamente 50 millones de filas; dim_taxi_zones, dim_payment_type, dim_vendor y dim_date están completamente cargadas; dbt test pasa de principio a fin; y analytics.taxi_zones_dict está activo y devuelve distritos mediante dictGet(). agg_hourly_zone_trips sigue vacía por diseño, no por un defecto, y permanecerá así hasta la transición del módulo 05.

Los paneles, el benchmark de ClickHouse frente a Snowflake y la transición corresponden al módulo 05, no a este.

El productor de Snowflake sigue en ejecución y el desfase entre Snowflake y ClickHouse continúa abierto. Nada de este módulo ha tocado el productor ni el script de migración. El módulo 05 cerrará ese desfase deliberadamente, mediante el mismo paso controlado de dos pasadas que anticipó el módulo 03. No detengas todavía el productor.

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