Snowflake MigrationClickHouse Workshops

Motores MergeTree

Cómo elegir un motor de la familia MergeTree y diseñar una clave ORDER BY que resulte provechosa.

ClickHouse almacena todos los datos en tablas respaldadas por una de las variantes del motor MergeTree. Si vienes de Snowflake, allí no existe un concepto equivalente: Snowflake toma internamente todas las decisiones de almacenamiento. En ClickHouse, elegir el motor adecuado es responsabilidad tuya, y una elección equivocada genera resultados incorrectos sin avisar.

Esta guía explica los motores que usarás en el laboratorio NYC Taxi y las trampas en las que tropiezan quienes migran desde Snowflake.


¿Qué es MergeTree?

MergeTree es el motor de almacenamiento principal de ClickHouse. Los datos se escriben en archivos columnares inmutables llamados partes. ClickHouse fusiona periódicamente las partes en segundo plano: las ordena, comprime y, si corresponde, transforma según las reglas del motor.

La consecuencia esencial es que una lectura puede ver varias versiones de una fila hasta que se produzca una fusión. La mayoría de los motores lo gestionan de forma transparente, pero algunos (especialmente ReplacingMergeTree) exigen entender el ciclo de vida de las fusiones para escribir consultas correctas.

Al crear una tabla MergeTree debes especificar ORDER BY. Esto determina:

  1. El orden físico de los datos dentro de cada parte
  2. El índice primario (disperso, por bloques y almacenado en memoria)
  3. En motores con deduplicación, las columnas que definen la «clave» de deduplicación

No existen conceptos separados de clave primaria, índice agrupado o clave de distribución. ORDER BY cumple todas esas funciones a la vez.


MergeTree

Úsalo cuando: la tabla solo recibe inserciones o las actualizaciones se gestionan externamente. No se necesita deduplicación.

CREATE TABLE default.some_events (
    event_id      String,
    occurred_at   DateTime64(3, 'UTC'),
    payload       String
)
ENGINE = MergeTree()
ORDER BY (occurred_at, event_id);

Características:

  • Las inserciones añaden datos como partes nuevas
  • No hay deduplicación: las filas duplicadas se conservan
  • Las fusiones optimizan el almacenamiento y la compresión, pero no alteran el contenido lógico
  • Las consultas leen todas las partes que coinciden con el intervalo del prefijo de ORDER BY

Cuándo falla: si insertas dos veces la misma fila (por ejemplo, al reintentar tras un fallo de red), ambas aparecen en los resultados. Esto es correcto para canalizaciones que solo insertan y donde no pueden existir duplicados. En toda tabla que reciba actualizaciones CDC o cargas reintentables, usa ReplacingMergeTree.


ReplacingMergeTree

Úsalo cuando: las filas pueden actualizarse (por ejemplo, correcciones de tarifas o cambios de estado) y quieres una sola fila por clave en los resultados.

CREATE TABLE analytics.fact_trips (
    trip_id       String,
    pickup_at     DateTime64(3, 'UTC'),
    fare_amount   Float64,
    updated_at    DateTime64(3, 'UTC'),
    -- ...
)
ENGINE = ReplacingMergeTree(updated_at)
ORDER BY (toStartOfMonth(pickup_at), pickup_at, trip_id);

Características:

  • Durante las fusiones en segundo plano, se deduplican las filas con la misma clave ORDER BY: solo se conserva la fila con el valor más alto en la columna de versión
  • La columna de versión (aquí updated_at) decide qué fila prevalece: mayor valor = más reciente = conservada
  • La deduplicación es asíncrona: hasta que se ejecuta una fusión, coexisten las versiones antigua y nueva

La trampa crítica: retraso de deduplicación

Entre fusiones, una consulta sin FINAL ve todas las versiones de una fila:

-- This may return multiple rows for the same trip_id
-- if the row has been updated since the last merge
SELECT * FROM analytics.fact_trips WHERE trip_id = 'abc123';

-- This returns exactly one row per trip_id, applying deduplication at query time
SELECT * FROM analytics.fact_trips FINAL WHERE trip_id = 'abc123';

FINAL fuerza la deduplicación durante la lectura. Es más lento que leer sin FINAL porque ClickHouse debe buscar claves duplicadas en todas las partes. En el laboratorio NYC Taxi, todas las consultas contra fact_trips usan FINAL.

Cuándo falla:

  • Omitir FINAL en una búsqueda puntual → devuelve silenciosamente filas duplicadas y las agregaciones cuentan de más
  • Usar una columna de versión incorrecta (que no aumenta al actualizar) → prevalecen valores antiguos
  • Usar MergeTree en vez de RMT para una tabla mutable → se acumulan todas las versiones y el recuento crece sin límite
  • Esperar deduplicación síncrona → un trabajo ETL lee justo después de insertar y ve duplicados

RMT con dbt: la estrategia incremental delete_insert elimina las filas del intervalo de claves del lote entrante antes de insertarlo, por lo que la tabla nunca tiene duplicados. Se sigue recomendando FINAL como medida de seguridad, aunque es menos crítico si la estrategia dbt es correcta.


AggregatingMergeTree

Úsalo cuando: la tabla almacena estados de agregación parciales que deben fusionarse durante las fusiones de fondo y combinarse durante la consulta.

CREATE TABLE analytics.agg_hourly_revenue (
    hour_bucket   DateTime,
    borough       String,
    fare_sum      AggregateFunction(sum, Float64),
    trip_count    AggregateFunction(count, UInt64)
)
ENGINE = AggregatingMergeTree()
ORDER BY (hour_bucket, borough);

Características:

  • Las filas con la misma clave ORDER BY se fusionan mediante la lógica combinadora de la función de agregación
  • Durante la consulta se usan combinadores con sufijo -Merge: sumMerge(fare_sum), countMerge(trip_count)
  • Suele alimentarlo una vista materializada que convierte inserciones sin procesar en estados parciales

Cuándo usarlo: AggregatingMergeTree sirve para datos preagregados cuyos estados parciales deben poder combinarse. En el laboratorio NYC Taxi, dbt reconstruye agg_hourly_zone_trips en cada ejecución: es una tabla de sustitución completa, no un acumulador de estados parciales. Usa ReplacingMergeTree para ella.

Cuándo falla: usar sum(fare_sum) en vez de sumMerge(fare_sum) durante la consulta trata el estado binario de agregación como un Float64 y devuelve números absurdos. Es un error silencioso de corrección.


CollapsingMergeTree

Úsalo cuando: necesitas eliminar o actualizar filas insertando una «fila de signo» (sign=1 para insertar, sign=-1 para cancelar). Es menos habitual, pero resulta útil para patrones CDC basados en eventos.

ENGINE = CollapsingMergeTree(sign)

Durante las fusiones, los pares de filas con sign=1 y sign=-1 para la misma clave se anulan. No se usa en el laboratorio NYC Taxi: ReplacingMergeTree con una columna de versión es más sencillo para el patrón de reintento de inserciones de esta carga.


MergeTree con TTL

Añade caducidad temporal de datos a cualquier variante MergeTree:

CREATE TABLE default.trips_raw (
    trip_id    String,
    pickup_at  DateTime64(3, 'UTC'),
    _synced_at DateTime DEFAULT now(),
    -- ...
)
ENGINE = ReplacingMergeTree(_synced_at)
ORDER BY (pickup_at, trip_id)
TTL toDate(pickup_at) + INTERVAL 2 YEAR;

El TTL se aplica durante las fusiones en segundo plano. Las filas caducadas se eliminan de las partes cuando estas se fusionan. El laboratorio no configura TTL: se conservan los cuatro años de datos. En producción, el TTL es esencial para controlar los costes de almacenamiento.


Elegir un motor: árbol de decisión

Does the table receive UPDATE or DELETE operations?
├── No (insert-only, e.g., event log, append-only stream)
│   └── MergeTree()
└── Yes
    ├── Do rows have a version/timestamp column that increases on update?
    │   ├── Yes → ReplacingMergeTree(version_col)
    │   └── No (full reload, e.g., dim tables rebuilt by dbt)
    │       └── MergeTree() — dbt atomic table swap (full rebuild) handles "upsert"
    └── Is the table a pre-aggregated accumulator with combinable states?
        └── AggregatingMergeTree()

Para el laboratorio NYC Taxi:

TablaMotorMotivo
trips_rawReplacingMergeTree(_synced_at)Los reintentos del script de migración y del productor después del cambio pueden escribir dos veces el mismo trip_id; _synced_at DEFAULT now() garantiza que prevalezca la escritura posterior
fact_tripsReplacingMergeTree(updated_at)Los viajes pueden corregirse; versión updated_at
agg_hourly_zone_tripsReplacingMergeTree(updated_at)Recálculo móvil = upsert; versión updated_at
Tablas dim_*MergeTreeRecarga completa mediante dbt; sin actualizaciones parciales
mv_hourly_revenueMV actualizableSe ejecuta según una programación y sustituye todo el resultado cada vez

Diseño de ORDER BY

La cláusula ORDER BY es la decisión de rendimiento más importante de una tabla de ClickHouse. Determina:

  1. Eficiencia del índice primario: las consultas que filtran por columnas del prefijo de ORDER BY omiten los bloques irrelevantes
  2. Tasa de compresión: los datos ordenados se comprimen mejor (los valores similares quedan juntos)
  3. Clave de deduplicación (para RMT/AMT): dos filas solo son duplicadas si coinciden sus columnas de ORDER BY

Reglas para diseñar ORDER BY:

  1. Coloca primero las columnas de cardinalidad baja (por ejemplo, borough, payment_type): más filas comparten valor y el índice puede omitir más bloques
  2. Coloca al final las columnas de cardinalidad alta (por ejemplo, trip_id, UUID): acotan el intervalo, pero comprimen peor al principio
  3. Deriva las columnas de los filtros reales de las consultas, no del esquema de origen
  4. En tablas RMT, la última columna debe ser el identificador único de fila (garantiza una fila por clave de negocio)

Antipatrón: copiar la clave primaria del origen como ORDER BY. Si TRIPS_RAW de Snowflake no tiene una ordenación explícita, copiar el orden del esquema de Snowflake (trip_id primero) da a ClickHouse un ORDER BY aleatorio, sin posibilidad de omitir bloques para ninguna consulta analítica.

Ejemplo de derivación para fact_trips:

Todas las consultas Q1–Q7 filtran de algún modo por pickup_at:

  • Q1: WHERE pickup_at >= ...
  • Q2: ORDER BY week, pickup_location_id
  • Q3: WHERE pickup_at >= CURRENT_DATE - 7
  • Q4: GROUP BY DATE_TRUNC('day', pickup_at)

Por tanto, pickup_at debe formar parte de ORDER BY y quedar cerca del principio. Usar toStartOfMonth(pickup_at) como primera columna crea un prefijo más grueso que permite podar a nivel de partición incluso sin una cláusula PARTITION BY. trip_id queda al final para aportar unicidad a RMT.

Resultado: ORDER BY (toStartOfMonth(pickup_at), pickup_at, trip_id)


PARTITION BY

PARTITION BY es opcional e independiente de ORDER BY. Crea particiones físicas de directorios; cada partición es un conjunto independiente de partes.

PARTITION BY toYYYYMM(pickup_at)

Usa PARTITION BY cuando:

  • Necesites eliminar eficientemente un intervalo temporal completo (ALTER TABLE DROP PARTITION '202401')
  • Quieras que el TTL opere por mes y no por fila
  • La tabla sea muy grande (>1 TB) y los metadatos por partición favorezcan la planificación de consultas

No uses PARTITION BY para sustituir ORDER BY. Un error habitual consiste en poner toYYYYMM(date) en PARTITION BY y omitirlo de ORDER BY, lo que impide omitir bloques dentro de una partición.

El laboratorio NYC Taxi no necesita PARTITION BY: el conjunto tiene 50 millones de filas (unos 8 GB comprimidos), muy dentro del intervalo de rendimiento de una sola partición.


Resumen de trampas principales

TrampaConsecuenciaSolución
Motor incorrecto para datos mutablesSe acumulan silenciosamente filas duplicadasUsa ReplacingMergeTree + FINAL
Falta FINAL en una consulta RMTLas agregaciones cuentan de más mientras se retrasan las fusionesAñade FINAL a todas las consultas analíticas sobre tablas RMT
ORDER BY tomado del esquema de origenConsultas lentas; no se omiten bloquesDeriva ORDER BY de los filtros reales de las consultas
Columna de cardinalidad alta al principio de ORDER BYMala selectividad del índiceCardinalidad baja primero y alta al final
Columna AggregateFunction consultada con sum() y no con sumMerge()Resultados numéricos absurdos sin avisarUsa siempre combinadores -Merge con AggregatingMergeTree
Columna de versión RMT que no aumenta monótonamentePrevalece una versión anterior de forma aleatoriaUsa una marca de tiempo que siempre se establezca en now() al actualizar

En esta página

ES