Elasticsearch MigrationClickHouse Workshops

Razonamiento de escenarios aplicados

Razonamiento detallado de ClickHouse detrás de los cinco escenarios aplicados de opción múltiple.

Completa primero las 20 preguntas aplicadas en la evaluación interactiva de la Parte 4. Los cinco antiguos escenarios escritos aparecen ahora como cuatro decisiones puntuadas cada uno. Esta referencia conserva ejemplos completos para debatir después.


Escenario 1 — Diseño del esquema

CREATE TABLE modelo

El siguiente CREATE TABLE materializa la decisión.

CREATE TABLE IF NOT EXISTS clickstream_events
(
    `event_id`   UUID                                                  CODEC(ZSTD(1)),
    `event_time` DateTime64(3)                                          CODEC(Delta, ZSTD(1)),
    `event_date` Date          DEFAULT toDate(event_time),
    `user_id`    UUID                                                  CODEC(ZSTD(1)),
    `session_id` UUID                                                  CODEC(ZSTD(1)),
    `product_id` String                                                CODEC(ZSTD(1)),
    `event_type` LowCardinality(String)                                CODEC(ZSTD(1)),
    `value_usd`  Decimal(18, 4)                                        CODEC(ZSTD(1)),
    `attributes` Map(LowCardinality(String), String)                   CODEC(ZSTD(1)),

    -- Skip indexes for non-prefix point-lookups
    INDEX idx_user_id    user_id    TYPE bloom_filter(0.01) GRANULARITY 4,
    INDEX idx_session_id session_id TYPE bloom_filter(0.01) GRANULARITY 4
)
ENGINE = MergeTree
PARTITION BY event_date
ORDER BY (event_type, product_id, event_time)
TTL event_date + INTERVAL 90 DAY DELETE
SETTINGS ttl_only_drop_parts = 1, index_granularity = 8192;

Justificación

  • Tipos:

    • event_id → UUID (16 bytes frente a 36 como String; rara vez se filtra).
    • user_id / session_id → UUID; cardinalidad ~50 M / ~500 M, por tanto no LowCardinality.
    • event_type → LowCardinality(String); 12 valores es el caso clásico. Enum serviría, pero LowCardinality(String) permite añadir un tipo 13 sin redesplegar.
    • product_id → String; ~2 M supera el rango normal de LowCardinality(String). Hay que medir, pero String es el valor predeterminado defendible.
    • attributes → Map(LowCardinality(String), String) para evolucionar con seguridad; promover claves usadas a MATERIALIZED.
    • value_usd → Decimal(18, 4) para evitar redondeo de ingresos; no hace falta Nullable.
  • ORDER BY (event_type, product_id, event_time):

    • Consulta nº 1 (embudo por producto) comienza con WHERE event_type IN ('view','add_to_cart','purchase') AND product_id = X; usa las dos columnas iniciales y poda muchos gránulos.
    • Consulta nº 3 (ingresos/minuto) filtra WHERE event_type = 'purchase', también prefijo.
    • Consulta nº 2 (historial de usuario) no usa el prefijo; para eso está bloom_filter en user_id.
    • Regla de cardinalidad: baja (12) → media (2 M) → alta/rango (timestamp). Columnas iniciales de baja cardinalidad mejoran la poda.
  • PARTITION BY event_date:

    • Una partición/día; con 5.000 millones de eventos/día evita miles de particiones pequeñas.
    • Coincide con TTL y permite que ttl_only_drop_parts = 1 elimine particiones completas.
    • No particiones por hora ni event_type; explotarían el recuento.
  • TTL event_date + INTERVAL 90 DAY DELETE con ttl_only_drop_parts = 1: elimina de una vez la partición diaria de 90 días, sin recorrer filas.

  • Índices de salto: bloom_filter en user_id (join de la Consulta nº 2) y session_id (uso futuro). No tokenbf_v1: son UUID, no texto. Granularidad 4 cubre 4 × 8192 = ~32K filas; ajusta midiendo latencia.

  • ¿AggregatingMergeTree para la Consulta nº 3? Sí, mediante una MV:

    CREATE MATERIALIZED VIEW revenue_by_minute_mv
    TO revenue_by_minute AS
    SELECT
        toStartOfMinute(event_time) AS minute,
        product_id,
        sumState(value_usd)         AS revenue,
        countState()                AS purchase_count
    FROM clickstream_events
    WHERE event_type = 'purchase'
    GROUP BY minute, product_id;

    Con 5.000 millones/día, consultar la tabla bruta bajo contención es lento; AggMergeTree desplaza el coste a la inserción. Las Consultas 1 y 2 tienen cardinalidad y variación excesivas para una MV.

Criterios de decisión

Las preguntas valoran:

  • MergeTree con ORDER BY cuyo prefijo se ajuste a la Consulta 1 o 3, no (event_id, event_time) ni (event_time, …)
  • Partición por fecha/semana, no event_type/hora
  • Map o JSON para attributes, no una columna por atributo
  • Al menos un bloom_filter en una columna de alta cardinalidad de la Consulta 2
  • TTL alineado con la partición

Una alternativa es errónea si coloca event_time primero en ORDER BY (instinto de ES perjudicial) o particiona por algo no derivado de fecha.


Escenario 2 — Plan de migración

Plan modelo de 6 meses

Fase 1 (semanas 1–2): descubrimiento y sandbox. Crea ClickHouse Cloud Production en la región de ES y un clúster paralelo de OTel Collector (3 nodos tras LB) que lea una derivación de la ingesta y al principio solo escriba en CH. Valida conexión, costes con 200K eventos/s sostenidos y contrapresión. Salida: CH ingiere 1 hora en doble flujo con desviación ≤1 % frente a ES.

Fase 2 (semanas 3–6): esquema + primeros 5 índices valiosos. Mapea campos de los 5 dashboards más consultados. Usa Map(LowCardinality(String), String) para los más de 60 campos grok y materializa los 10 que aparecen en ≥3 dashboards. Define MV hacia cualquier AggregatingMergeTree. Salida: las consultas principales de los 20 dashboards tienen equivalente CH igual o más rápido.

Fase 3 (semanas 7–12): ejecución paralela + dashboards. Activa escritura doble. Migra ~50 búsquedas guardadas/semana, primero las que exploran >1 TB. Ejecuta el script de validación dos veces al día e investiga desviaciones >5 %. Salida: 300 búsquedas funcionan; p95 de CH ≤ p95 de ES durante 7 días.

Fase 4 (semanas 13–18): retirar Logstash. Sustituye Logstash por OTel Collector. Los más de 60 grok se convierten en procesadores transform y columnas MATERIALIZED. Ejecuta ambos 2 semanas y elimina Logstash tras confirmar paridad. Salida: 7 días sin tráfico por Logstash; alertas en HyperDX Alerts o equivalente.

Fase 5 (semanas 19–22): transición. Detén Filebeat → Logstash, elimina Logstash y detén nueva ingesta ES, conservando el clúster solo lectura. Supervisa costes y latencia. Salida: ES solo lectura 14 días sin incidentes ni reversión.

Fase 6 (semanas 23–26): desactivación y conformidad. Snapshot final completo de ES a S3 en la fecha de transición, desactiva el clúster y documenta el bucket. Cambia el modelo mental de «ES es la verdad» a «ES es el retrovisor de 90 días». Salida: clúster destruido, snapshot validado y runbook actualizado.


Estrategia para más de 60 campos grok

  • Mes 1: los 60 llegan a LogAttributes Map(LowCardinality(String), String).
  • Mes 2: identifica los 10 más referenciados y promueve con ALTER TABLE … ADD COLUMN field MATERIALIZED LogAttributes['field'].
  • Mes 4: repite la consulta y promueve nuevos campos frecuentes; la cola larga queda en Map.
  • Mes 6: auditoría final. El esquema suele converger en 15–25 columnas + Map.

Así se evita la trampa de «diseñar de antemano el esquema perfecto».


Decisión sobre Logstash: sustituir por OTel Collector

Los más de 60 grok dominan computación y complejidad. transform (OTTL) cubre ~80 %; el resto se convierte en columnas MATERIALIZED, trasladando el análisis al almacenamiento y calculándolo una vez al insertar.

Vector es viable y VRL resulta familiar a usuarios de Logstash, pero añade otra herramienta. OTel Collector mantiene coherencia.

No conserves Logstash. La historia es una herramienta menos, no «ES + Logstash + Kibana → CH + Logstash + HyperDX».


Retención de conformidad (capa fría de 1 año)

CH Cloud no tiene ILM de snapshots S3 integrado, pero funcionan dos patrones:

  1. Exportación diaria al S3 del cliente. Programa INSERT INTO FUNCTION s3('s3://archive/year/month/day.parquet') SELECT * FROM otel_logs_v2 WHERE event_date = today() - 1. El bucket conserva un año de Parquet, consultable con Athena, Trino o s3() desde CH temporal. S3 IA cuesta ~$23/TB/mes, mucho menos que un año activo.
  2. Almacenamiento por capas de ClickHouse Cloud, donde esté disponible. Después de 90 días puede costar ~70 % menos, pero auditarlo es más complejo porque los datos siguen dentro de CH.

Recomendamos (1): el auditor quiere que «logs de agosto de 2025» sea una consulta S3 autoservicio, no un proyecto para iniciar un clúster.


Reversión si falla tres semanas después

Se presupone escritura doble y algunos dashboards migrados:

  1. Detener escrituras solo CH: vuelve el collector a escritura única en ES.
  2. Volver a Kibana: las búsquedas originales siguen presentes.
  3. Mantener CH solo lectura para analizar los datos de la ventana paralela.
  4. No eliminar artefactos migrados: conserva esquema y configuraciones; al corregir la causa, reinicia en Fase 4, no repitas 1–3.

Propiedad esencial: en cualquiera de las primeras 5 fases se vuelve a ES cambiando una configuración de OTel Collector.


3 riesgos principales y mitigaciones

  1. Esquema fijado demasiado pronto. Programa revisiones de campos frecuentes en meses 2, 4 y 6; las promociones ALTER TABLE no son destructivas.
  2. Cinco alertas se comportan distinto. Crea MV según el módulo 03, Paso 8 Opción B en semana 8 y ejecútalas en sombra 4 semanas; compara disparos con Kibana.
  3. Patrones grok con regex no admitidas. Enumera 60 en semana 1, identifica lookbehind/condicionales y reescribe como OTTL o columnas regexpExtract. No lo descubras en semana 14.

Criterios de decisión

Se valoran planes que:

  • Tienen fases y criterios de salida claros
  • Sustituyen Logstash por OTel Collector o Vector; conservarlo es un fallo
  • Usan Map + promoción selectiva, no 60 columnas iniciales ni JSON para todo
  • Nombran un patrón de conformidad concreto (S3 Parquet, capas o equivalente)
  • Definen reversión sin reimportar CSV

Escenario 3 — Depuración A (desviación de latencia p95)

Respuesta modelo

Tres causas probables:

Causa 1: unidades diferentes. ES solía guardar milisegundos (response_time en segundos ×1000 o ya ms). Duration de trazas en ClickHouse está en nanosegundos; quantile(0.95)(Duration) sin dividir por 1e6 da ~1000× más. También LogAttributes['run_time'] podría estar en segundos mientras ES tenía run_time_ms convertido. Comprobación: Ejecuta SELECT min(Duration), max(Duration), avg(Duration) FROM otel_traces y compara min/max/avg de ES. Distinta magnitud = unidades.

Causa 2: algoritmos de cuantiles diferentes. percentiles de Elasticsearch usa HDR Histogram o T-Digest, ambos aproximados. quantile() de ClickHouse también es aproximado, pero distinto; colas largas pueden divergir 10–25 %. quantileExact() es exacto y lento; quantileTDigest() se acerca a ES. Comprobación: Ejecuta quantileExact(0.95)(Duration). Si el exacto cae entre ambas aproximaciones, solo difieren las estrategias.

Causa 3: ventanas o muestras distintas por reloj/filtros. now() - INTERVAL 1 HOUR puede representar otro instante, o TimestampTime tiene precisión de segundos mientras @timestamp de ES usa milisegundos. Cinco segundos pueden mover p95 20 %. El tipo materializado en este caso es DateTime. Comprobación: Fija BETWEEN '2026-05-09 04:00:00' AND '2026-05-09 05:00:00' en ambos. Si convergen, era alineación temporal.

Criterios de decisión

Se valora:

  • Una causa de unidades/tipos (Duration ns frente a ms)
  • Algoritmo de cuantiles O alineación temporal
  • Una comprobación concreta por causa, no solo «investigar»

Escenario 4 — Depuración B (búsqueda lenta por clave Map)

Diagnóstico

LogAttributes['request_id'] es una desreferencia de Map por fila; ClickHouse busca 'request_id' dentro del Map. No hay índice de salto aplicable y debe materializar el Map en cada fila del rango. TraceId, en cambio, es String superior con bloom_filter en el esquema de otel_logs_v2, que elimina ~99,9 % de gránulos antes de leer. La columna Map implicada es LogAttributes.

La ralentización 200× no significa que Map sea lento; la optimización del índice no se aplica a claves.

Corrección rápida: índice de Bloom en mapValues

ALTER TABLE otel_logs_v2 ADD INDEX idx_request_id
    mapValues(LogAttributes)
    TYPE bloom_filter(0.01)
    GRANULARITY 4;

ALTER TABLE otel_logs_v2 MATERIALIZE INDEX idx_request_id;

SELECT *
FROM otel_logs_v2
WHERE indexHint(has(mapValues(LogAttributes), 'abc123'))
  AND LogAttributes['request_id'] = 'abc123';

Crea un Bloom por gránulo de todos los valores. La pista has(mapValues(...)) poda candidatos y el predicado con clave garantiza la corrección. Contrapartida: cubre todos los valores, con más colisiones que un Bloom dedicado. Mide con EXPLAIN indexes = 1; una DDL correcta no demuestra poda útil.

Corrección duradera: promover request_id a columna materializada

ALTER TABLE otel_logs_v2 ADD COLUMN request_id String MATERIALIZED LogAttributes['request_id'];
ALTER TABLE otel_logs_v2 ADD INDEX idx_request_id request_id TYPE bloom_filter(0.01) GRANULARITY 4;
ALTER TABLE otel_logs_v2 MATERIALIZE COLUMN request_id, INDEX idx_request_id;

Ahora request_id tiene compresión y Bloom propios. Consulta WHERE request_id = 'abc123' en lugar de WHERE LogAttributes['request_id'] = 'abc123'. Cada promoción exige reescribir una columna una vez; compensa para claves usadas, no para la cola larga.

Es el patrón de la Decisión 3 del ADR de la Parte 2: Map por defecto, promover claves usadas.

Criterios de decisión

Se valora:

  • Identificar «sin vía de índice de salto en la desreferencia», no «ClickHouse es lento»
  • Bloom de valores Map + predicado/pista compatible como corrección rápida medida
  • Promoción a columna como solución duradera
  • Reconocer colisiones más amplias frente a coste ALTER

Escenario 5 — Análisis de contrapartidas (diseños de alertas)

Tabla comparativa — respuesta modelo

PropiedadDiseño A (en consulta)Diseño B (MV preagregada)
Lectura por evaluaciónAlta: cada 60 s explora 5 minutos. 1 M eventos/s × 300 s = 300 M filas.Baja: ~1 fila/minuto/servicio; la consulta lee ~5 filas en microsegundos.
Escritura por filaCeroPequeña pero no nula: cada inserción actualiza el bucket; normalmente menos del 5 %.
Retraso~60 s + latencia~60 s; la MV se actualiza síncronamente al insertar
Escala de cardinalidad (100 servicios × 50 endpoints)Mismo coste: tabla bruta5000 grupos × 1/min = 7,2 M filas/día, diminuto frente a datos brutos

Recomendación a 1 M eventos/s — gana el Diseño B

El equilibrio llega cuando explorar datos brutos supera el coste de actualizar la MV:

  • A explora 300 M filas por evaluación. Incluso a 1 GB/s por CPU y caché fría son segundos, cada minuto, millones de segundos de CPU por semana.
  • B amortiza la agregación en la inserción, pagándola una vez por fila en vez de cada minuto. Los lotes columnares hacen casi gratuita la actualización del bucket.

Intuición: volver a explorar datos brutos es el patrón de transformaciones de Elasticsearch; MV preagregadas son el patrón nativo de ClickHouse. A ~100K eventos/s pasan de conveniencia a requisito; a 1 M, A gastaría más CPU que la carga supervisada.

Nota: cuándo sigue teniendo sentido el Diseño A

  • Entornos de menos de 10K eventos/s
  • Alertas ad hoc o temporales sin querer aprovisionar una MV
  • Consultas que cambian cada semana; el esquema de MV es fijo

Criterios de decisión

Se valora:

  • Coste de B como O(buckets) y A como O(eventos brutos explorados)
  • Contrapartida: A nada al escribir y mucho al leer; B algo al escribir y casi nada al leer
  • Elegir B a 1 M eventos/s por amortización, no solo «MV es más rápida»
  • Reconocer agregación al insertar frente a al consultar, el patrón de AggregatingMergeTree logs_summary_1min de la Parte 3

Revisión de la puntuación aplicada

La evaluación exige 16 de 20 respuestas aplicadas. Si quedaste por debajo, usa los criterios para localizar el escenario:

  • Ejercicio 1: tratar ClickHouse como Elasticsearch (una columna por atributo, particionar por event_type u ordenar primero por timestamp). Revisa schema.sql y la hoja de la Parte 2.
  • Ejercicio 2: olvidar reversión o conformidad. La migración es reversible hasta desactivar ES.
  • Ejercicio 4: creer que Map es inherentemente lento; el problema es qué optimizaciones se aplican.

Vuelve a la Parte 4 y repite el componente aplicado tras revisar el escenario fallado.

En esta página

ES