Elasticsearch MigrationClickHouse Workshops

Solución de la hoja del modelo de datos

Respuestas modelo para mapear el modelo de datos de Elasticsearch a ClickHouse.

Nota: Los recuentos, tamaños y nombres variarán; estos valores proceden de una ejecución de referencia tras ~5 días. Importa la forma de las respuestas: tipos, cardinalidades, finalidad del enriquecimiento y justificación del diseño.


1. Flujo de datos: logs-web_access-lab

Estado actual

  • Documentos: ~31 millones
  • Índices subyacentes: 3 (generación 3; uno por día o 5 GB)
  • Tamaño primario: ~8,5 GB
  • Shards por índice: 2 primarios (plantilla lab-logs-settings)
  • Réplicas: 1, UNASSIGNED en el laboratorio de un nodo
  • Campos únicos: ~69 de usuario, sin metacampos como _id, _index; incluye geo.*, user_agent_parsed.* y .keyword dinámicos
  • Alta cardinalidad: remote_addr (~65 mil), request_path (~200, Zipf moderada), trace.id si se une con APM, user_agent (cientos)
  • Campos más consultados: @timestamp, status, request_path, request_type, service, geo.country_name, user_agent_parsed.name, event.severity, run_time
  • Política: lab-observability-policy
  • Fases: hot (prioridad 100 y rollover con max_age=1d o max_size=5gb), warm (2 d: shrink a 1 shard, forcemerge a 1 segmento, prioridad 50), delete (30 d: eliminar + snapshot)
  • Rollover: max_size=5 GB O max_age=1 d
  • Eliminación: 30 días
  • Pipelines: default-enrichment → web-access-enrichment
  • Procesadores:
    • set event.ingested = _ingest.timestamp (default-enrichment) — marca de ingesta
    • geoip on remote_addr → geo.* — country_name, city_name y geopunto location
    • user_agent on user_agent → user_agent_parsed.* — name, version, os, device
    • set event.severity = "info" — severidad predeterminada
    • script (Painless) — warn en 4xx y error en 5xx
  • Campos enriquecidos: geo.country_name, geo.city_name, geo.location, user_agent_parsed.name/version/os.*/device.name, event.severity, event.ingested

Diseño propuesto de tabla

CREATE TABLE otel_logs_web_access
(
    Timestamp       DateTime64(9) CODEC(Delta, ZSTD),
    ServiceName     LowCardinality(String),
    SeverityText    LowCardinality(String),
    Body            String CODEC(ZSTD(3)),
    LogAttributes   Map(LowCardinality(String), String),
    -- materialized columns for hot-path dashboard fields
    RemoteAddr      IPv4      MATERIALIZED toIPv4OrDefault(LogAttributes['remote_addr']),
    Status          UInt16    MATERIALIZED toUInt16OrZero(LogAttributes['status']),
    RequestPath     String    MATERIALIZED LogAttributes['request_path'],
    RequestMethod   LowCardinality(String) MATERIALIZED LogAttributes['request_type'],
    RunTime         Float32   MATERIALIZED toFloat32OrZero(LogAttributes['run_time']),
    CountryName     LowCardinality(String) MATERIALIZED dictGetOrDefault('geo_ip_dict', 'country_name', RemoteAddr, ''),
    UserAgentName   LowCardinality(String) MATERIALIZED extract(LogAttributes['user_agent'], '^([A-Za-z]+)')
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(Timestamp)        -- 1 month per partition; matches retention granularity
ORDER BY (ServiceName, Status, Timestamp)  -- dashboards filter by service + status
TTL toDateTime(Timestamp) + INTERVAL 30 DAY DELETE;
  • Partición: toYYYYMM(Timestamp), una por mes. Una granularidad más fina (toDate) crea miles de partes; coincide con la retención de 30 días.
  • ORDER BY: (ServiceName, Status, Timestamp) porque los dashboards filtran primero ServiceName y Status; Timestamp permite podar rangos dentro de cada bloque.
  • Columnas frente a Map: materializa campos frecuentes (Status, RequestPath, RunTime) y conserva campos fríos (referer, size, user_agent_parsed.device.name) en LogAttributes. Se calculan al insertar y se almacenan como columnas reales.
  • GeoIP: diccionario (diseño IP_TRIE sobre CSV MaxMind GeoLite2) + dictGet() materializado. Los agentes permanecen sin estado.
  • User-agent: procesador user_agent de OTel o columna con extract() / regex. El procesador es más limpio y se mapea directamente desde Elastic.
  • El campo original user_agent puede seguir en el Map para diagnóstico.
  • TTL: Timestamp + INTERVAL 30 DAY DELETE, igual que ILM.
  • Fases ILM innecesarias: rollover, shrink, forcemerge, set_priority y migración de capas; Cloud usa objetos con caché automática.

2. Flujo de datos: logs-application-lab

Estado actual

  • Documentos: ~12 millones
  • Índices: 2
  • Tamaño primario: ~3,9 GB
  • Shards/réplicas: 2 primarios / 1 réplica (no asignada)
  • Campos: ~42
  • Alta cardinalidad: trace_id, span_id, message, error.stack
  • Consultas comunes: @timestamp, level, service, event.severity, trace_id, message
  • Política: lab-observability-policy, igual que web
  • Fases/rollover/eliminación: iguales (hot 5 GB/1 d, warm 2 d, delete 30 d)
  • Pipelines: default-enrichment → app-log-enrichment
  • Procesadores:
    • set event.ingested
    • set event.severity = {{level}} — copiar level
    • lowercase event.severity
    • dissect on message (rara vez coincide; deja _tmp.*)
    • remove _tmp* — limpieza
  • Campos: event.severity, event.ingested

Diseño propuesto de tabla

CREATE TABLE otel_logs_application
(
    Timestamp      DateTime64(9) CODEC(Delta, ZSTD),
    ServiceName    LowCardinality(String),
    SeverityText   LowCardinality(String),
    Body           String CODEC(ZSTD(3)),
    TraceId        String,
    SpanId         String,
    LogAttributes  Map(LowCardinality(String), String),
    -- skip index for trace-ID needle-in-haystack lookup
    INDEX trace_id_bf TraceId TYPE bloom_filter(0.01) GRANULARITY 4,
    -- full-text skip index for `Body` (text index preferred >= 26.2; tokenbf_v1 deprecated)
    INDEX body_tokens Body TYPE text GRANULARITY 4
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(Timestamp)
ORDER BY (ServiceName, SeverityText, Timestamp)
TTL toDateTime(Timestamp) + INTERVAL 30 DAY DELETE;
  • Partición: toYYYYMM(Timestamp), misma razón.
  • ORDER BY: (ServiceName, SeverityText, Timestamp) para «errores del servicio X en los últimos 10 minutos».
  • TraceId / SpanId: columnas superiores con bloom_filter en TraceId. NO en la clave primaria, lo que destruiría la localidad temporal.
  • Severidad: MATERIALIZED lowerUTF8(LogAttributes['level']) reproduciría lowercase; aquí se almacena como SeverityText desde severityparser de OTel.
  • dissect: omitirlo; si hace falta análisis estructurado, usar regex_parser de OTel, no una MV.
  • TTL: 30 días.
  • Fases ILM: todas innecesarias salvo delete.

3. Flujo de datos: logs-infrastructure-lab

Estado actual

  • Documentos: ~18,6 millones
  • Índices: 2
  • Tamaño primario: ~3,7 GB
  • Shards/réplicas: 2 primarios / 1 réplica
  • Campos: ~35
  • Alta cardinalidad: message, log_message, pid
  • Baja cardinalidad: hostname (~10 valores, k8s-node-01..10), process (~10)
  • Consultas comunes: @timestamp, hostname, process, event.severity, log_message
  • Política/fases/eliminación: iguales (lab-observability-policy)
  • Pipelines: default-enrichment → infra-log-parsing
  • Procesadores:
    • set event.ingested
    • grok %{SYSLOGTIMESTAMP}%{HOSTNAME}%{WORD:process}[...] — extrae 4 campos
    • set event.severity = "info"
    • script (Painless) — eleva a warn/error según palabras de log_message
  • Campos: syslog_timestamp, hostname, process, pid, log_message, event.severity, event.ingested

Diseño propuesto de tabla

CREATE TABLE otel_logs_infrastructure
(
    Timestamp      DateTime64(9) CODEC(Delta, ZSTD),
    Hostname       LowCardinality(String),
    Process        LowCardinality(String),
    Pid            UInt32,
    SeverityText   LowCardinality(String),
    Body           String CODEC(ZSTD(3)),                  -- raw syslog line
    LogMessage     String CODEC(ZSTD(3)),                   -- extracted message body
    LogAttributes  Map(LowCardinality(String), String),
    INDEX body_tokens LogMessage TYPE text GRANULARITY 4  -- text index preferred >= 26.2; tokenbf_v1 deprecated
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(Timestamp)
ORDER BY (Hostname, Process, Timestamp)
TTL toDateTime(Timestamp) + INTERVAL 30 DAY DELETE;
  • Partición/ORDER BY: Los dashboards filtran por Hostname → Process → tiempo; ambos comprimen bien y permiten poda temporal.
  • Sustituir grok: ejecutar regex_parser en OTel Collector. Una columna MATERIALIZED con extractAllGroupsVertical() pagaría CPU en cada inserción de ClickHouse.
  • Severidad: MATERIALIZED multiIf(positionCaseInsensitive(LogMessage, 'error') > 0, 'error', positionCaseInsensitive(LogMessage, 'warn') > 0, 'warn', 'info') reproduce Painless.
  • TTL: 30 días.
  • Fases ILM: solo hace falta delete.

4. Flujos de datos: traces-apm-* y logs-apm.* (Carga 2 — OTel Demo)

4a. Trazas: traces-apm-*

  • Documentos: ~13,8 millones el día 1, creciendo ~14 millones/día (1 índice/día)
  • Índices: 2–3, según duración
  • Campos hoja: ~200 (recursos/atributos OTel + control APM)
  • service.name distintos: 19 (16 microservicios + frontend-proxy, frontend-web, sample-order-app)
  • processor.event: ~50 % transaction, ~50 % span
  • Lenguajes (service.language.name, 10, incluido unknown de Envoy): nodejs, cpp, python, dotnet, rust, java, php, ruby, go, unknown
  • Cardinalidad transaction.name: ~7 000
  • Percentiles transaction.duration.us: p50 ≈ 2,5 ms · p95 ≈ 75 ms · p99 ≈ 1,4 s
  • event.outcome: ~99,5 % success · ~0,5 % failure (off predeterminado en flagd)

4b. Logs: logs-apm.app.* y logs-apm.error-default

  • El patrón logs-apm.* reúne todos estos flujos de aplicaciones y errores.
  • Flujos logs-apm.*: 19 (18 logs-apm.app.<service>-default + 1 logs-apm.error-default)
  • Documentos totales: decenas de millones; frontend_web solo llega a ~28 millones/día
  • 3 servicios principales: frontend_web (~27,8 M) · frontend_proxy (~2,0 M) · product_catalog (~1,1 M); después cart (~1 M), currency (~400 k), recommendation (~300 k)
  • Observación: frontend_web domina 10× porque Next.js registra cada renderizado. El volumen por servicio está muy sesgado.

4c. Destino propuesto: otel_traces

Usa el esquema predeterminado de clickhouseexporter de OTel Collector. Se mapea directamente desde OTLP y gestiona lotes y migraciones.

-- What the exporter creates for you (abridged):
CREATE TABLE otel_traces
(
    Timestamp             DateTime64(9) CODEC(Delta, ZSTD(1)),
    TraceId               String CODEC(ZSTD(1)),
    SpanId                String CODEC(ZSTD(1)),
    ParentSpanId          String CODEC(ZSTD(1)),
    TraceState            String CODEC(ZSTD(1)),
    SpanName              LowCardinality(String) CODEC(ZSTD(1)),
    SpanKind              LowCardinality(String) CODEC(ZSTD(1)),
    ServiceName           LowCardinality(String) CODEC(ZSTD(1)),
    ResourceAttributes    Map(LowCardinality(String), String) CODEC(ZSTD(1)),
    SpanAttributes        Map(LowCardinality(String), String) CODEC(ZSTD(1)),
    Duration              Int64 CODEC(ZSTD(1)),
    StatusCode            LowCardinality(String) CODEC(ZSTD(1)),
    StatusMessage         String CODEC(ZSTD(1)),
    Events.Name           Array(String) CODEC(ZSTD(1)),
    Events.Timestamp      Array(DateTime64(9)) CODEC(Delta, ZSTD(1)),
    Events.Attributes     Array(Map(LowCardinality(String), String)) CODEC(ZSTD(1)),
    Links.TraceId         Array(String) CODEC(ZSTD(1)),
    Links.SpanId          Array(String) CODEC(ZSTD(1)),
    Links.TraceState      Array(String) CODEC(ZSTD(1)),
    Links.Attributes      Array(Map(LowCardinality(String), String)) CODEC(ZSTD(1)),
    INDEX idx_trace_id TraceId TYPE bloom_filter(0.001) GRANULARITY 1  -- ← we add this
)
ENGINE = MergeTree
PARTITION BY toDate(Timestamp)
ORDER BY (ServiceName, SpanName, Timestamp)
TTL toDateTime(Timestamp) + INTERVAL 30 DAY DELETE;
  • ORDER BY: (ServiceName, SpanName, Timestamp), para latencia de frontend.checkout en 15 minutos; predeterminado del exportador.
  • Paso obligatorio: el exportador no trae índice de salto en TraceId, por lo que la Consulta 5 exploraría todo. Ejecuta:
    ALTER TABLE otel_traces ADD INDEX trace_id_bf TraceId TYPE bloom_filter(0.01) GRANULARITY 4;
    ALTER TABLE otel_traces MATERIALIZE INDEX trace_id_bf;
  • TTL: 30 días.

4d. Destino propuesto: otel_logs (logs APM)

Una tabla otel_logs, no 19.

El destino consolidado sigue llamándose otel_logs.

Elasticsearch divide por servicio para limitar mapeos por índice y lograr ILM granular. ClickHouse invierte la relación; una tabla:

  • Comprime ServiceName como LowCardinality(String) en ~1 byte/fila.
  • Permite podar por servicio con ORDER BY (ServiceName, ...).
  • Evita la sobrecarga de partes/merges de 19 tablas.
  • Permite JOIN otel_logs USING (TraceId) otel_traces para correlación.
CREATE TABLE otel_logs
(
    Timestamp       DateTime64(9) CODEC(Delta, ZSTD),
    TraceId         String,
    SpanId          String,
    SeverityText    LowCardinality(String),
    SeverityNumber  Int32,
    ServiceName     LowCardinality(String),
    Body            String CODEC(ZSTD(3)),
    ResourceAttributes Map(LowCardinality(String), String),
    LogAttributes      Map(LowCardinality(String), String),
    INDEX idx_trace_id TraceId TYPE bloom_filter(0.01) GRANULARITY 4
)
ENGINE = MergeTree
PARTITION BY toYYYYMM(Timestamp)
ORDER BY (ServiceName, SeverityText, Timestamp)
TTL toDateTime(Timestamp) + INTERVAL 30 DAY DELETE;
  • ORDER BY (ServiceName, SeverityText, Timestamp) cubre errores de un servicio en N minutos.
  • Correlación TraceId mediante columna superior + Bloom; logs y trazas comparten valor y nombre.
  • Mantener TraceId fuera de la clave preserva la localidad temporal.
  • TTL: 30 días.

5. Referencia de latencia de consultas

Mediana de 3 ejecuciones contra ES en localhost:9200.

Consulta.took de ES (mediana, ms)Observaciones
Q1 — 10 rutas principales (estado 200)1Primera ejecución fría 200–1 500 ms; calientes ~1 ms. Terms sobre request_path.keyword.
Q2 — Recuento 5xx/min en 1 h4Ventana pequeña y cardinalidad baja.
Q3 — Búsqueda por trace.id1Fría 100–200 ms; caliente ~1 ms. Búsqueda puntual en índice invertido.

(Valores fríos habituales: Q1=254 ms · Q2=6 ms · Q3=158 ms. Ejecuta dos veces antes de registrar el estado caliente.)

Nota didáctica: Importa el patrón: Q3 es casi gratis en ES porque todo trace.id está en un índice invertido. En ClickHouse se consigue con bloom_filter (Decisión 3 del ADR). Sin él, WHERE TraceId = ? explora otel_traces completo.

La Parte 3 repite las consultas. Q1 y Q2 deberían ser más rápidas por la exploración columnar de Status; Q3 debe acercarse a ES con el Bloom. Si tarda más de 10×, probablemente el índice no se materializó en gránulos históricos; revisa la Sección 4c.


6. Observaciones generales

  • Almacenamiento total en los 3 flujos: ~16 GB primarios (web 8,5 + app 3,9 + infra 3,7). ClickHouse suele comprimir 10–15× mejor; espera 1–2 GB.
  • Procesadores que pueden vivir en cliente: geoip, user_agent, grok, set event.ingested y severidad, todos en OTel Collector (transform, geoip, user_agent, regex_parser, attributes). En el borde mantienen el backend mínimo; el coste de CPU escala con la flota, frente a dictGet() central.
  • Veredicto ILM:
    • rollover: innecesario, una tabla con particiones.
    • shrink: innecesario, shards lógicos autoescalables.
    • forcemerge: innecesario, merges automáticos.
    • set_priority: innecesario, sin prioridad por capas.
    • migrate (hot→warm→cold): innecesario, objetos con caché local.
    • delete a 30 d: necesario, mediante TTL.
  • Limitar almacenamiento al 30 % de ES: elimina size y referer, elimina Body de stack traces tras 7 días con TTL … RECOMPRESS o TTL por columna (TTL Timestamp + INTERVAL 7 DAY DELETE WHERE SeverityText = 'info') y reduce @timestamp a DateTime (4 bytes) en lugar de DateTime64(9) (8 bytes) cuando no haga falta precisión inferior al segundo.

En esta página

ES