Elasticsearch MigrationClickHouse Workshops

Solução da planilha de modelo de dados

Respostas-modelo para a planilha de mapeamento do modelo de dados do Elasticsearch para o ClickHouse.

Observação: As contagens de documentos, os tamanhos e os nomes dos índices serão diferentes no seu ambiente — os valores abaixo são de uma execução de referência após ~5 dias de geração de dados. O importante é o formato das respostas: tipos, cardinalidades, finalidade do enriquecimento e justificativa do projeto das tabelas.


1. Fluxo de dados: logs-web_access-lab

Estado atual

  • Total de documentos: ~31 milhões
  • Número de índices subjacentes: 3 (geração 3, um rotacionado por dia ou 5 GB)
  • Tamanho total dos shards primários: ~8,5 GB (nos 3 índices subjacentes)
  • Shards por índice subjacente: 2 primários (do modelo de componente lab-logs-settings)
  • Réplicas por índice subjacente: 1 réplica configurada — UNASSIGNED no laboratório de nó único (esperado; não há um segundo nó para hospedá-la)
  • Campos únicos no mapeamento: ~69 campos de usuário (sem contar metacampos do ES como _id, _index). Inclui subcampos geo.*, user_agent_parsed.* adicionados dinamicamente e subcampos .keyword.
  • Campos de alta cardinalidade: remote_addr (~65 mil), request_path (~200 com distribuição Zipf — moderada), trace.id (se associado ao APM), user_agent (várias centenas)
  • Campos usados na maioria das consultas (dos dashboards da Parte 1): @timestamp, status, request_path, request_type, service, geo.country_name, user_agent_parsed.name, event.severity, run_time
  • Política ILM: lab-observability-policy
  • Fases do ILM: hot (prioridade 100, rollover em max_age=1d ou max_size=5gb), warm (aos 2 d — shrink para 1 shard, forcemerge para 1 segmento, prioridade 50), delete (aos 30 d — exclusão + exclusão de snapshot)
  • Condição de rollover: max_size=5 GB OU max_age=1 d
  • Idade para exclusão: 30 dias
  • Pipelines de ingestão: default-enrichment (padrão, definido em todos os fluxos de dados) → web-access-enrichment (aplicado por roteamento condicional do Filebeat)
  • Processadores do pipeline:
    • set event.ingested = _ingest.timestamp (default-enrichment) — registrar o momento da ingestão
    • geoip on remote_addr → geo.* — adicionar country_name, city_name e geoponto location
    • user_agent on user_agent → user_agent_parsed.* — analisar name, version, os e device
    • set event.severity = "info" — severidade padrão
    • script (Painless) — substituir a severidade por warn em 4xx e error em 5xx
  • Campos enriquecidos: geo.country_name, geo.city_name, geo.location, user_agent_parsed.name/version/os.*/device.name, event.severity, event.ingested

Projeto proposto da tabela do ClickHouse

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;
  • Chave de partição: toYYYYMM(Timestamp) — uma partição por mês. Granularidade adequada para exclusões eficientes; granularidade excessiva (toDate) cria milhares de partes. Corresponde à retenção de 30 dias.
  • ORDER BY: (ServiceName, Status, Timestamp) — consultas de dashboard (taxa de solicitações por serviço, contagem de 5xx) filtram primeiro por ServiceName e depois por Status. Timestamp por último permite podar intervalos de tempo em cada bloco ordenado por serviço e status.
  • Colunas de nível superior vs. Map: Coloque os campos usados em consultas frequentes de dashboard como colunas materializadas (Status, RequestPath, RunTime); mantenha campos pouco acessados/de diagnóstico (referer, size, user_agent_parsed.device.name) em LogAttributes. Colunas materializadas são calculadas na inserção e armazenadas como colunas reais, proporcionando o desempenho completo da varredura colunar.
  • Enriquecimento GeoIP: Dicionário (layout IP_TRIE sobre CSV MaxMind GeoLite2) + dictGet() em uma coluna materializada. Transfira o enriquecimento ao mecanismo de armazenamento para manter os agentes de ingestão sem estado e evitar a reimplantação dos coletores após mudanças de esquema.
  • Análise de user-agent: Use (a) o processador user_agent do OTel (simples: emite name/os/device analisados como atributos) ou (b) uma coluna materializada com extract() / regex. O processador do OTel é mais adequado à migração porque o processador user_agent da Elastic já existe e corresponde diretamente a ele.
  • TTL: Timestamp + INTERVAL 30 DAY DELETE. Corresponde à idade de exclusão do ILM.
  • Fases do ILM desnecessárias no CH Cloud: rollover (nenhum índice para rotacionar), shrink (shards lógicos escalam automaticamente), forcemerge (mesclagens em segundo plano fazem isso automaticamente), set_priority (nenhum cache em camadas exposto), migração de camada (hot/warm/cold é irrelevante — todos os dados ficam em armazenamento de objetos com cache automático).

2. Fluxo de dados: logs-application-lab

Estado atual

  • Total de documentos: ~12 milhões
  • Índices subjacentes: 2
  • Tamanho dos shards primários: ~3,9 GB
  • Shards/réplicas: 2 primários / 1 réplica (réplica não atribuída — nó único)
  • Campos do mapeamento: ~42 campos de usuário
  • Campos de alta cardinalidade: trace_id, span_id, message, error.stack (quando presente)
  • Campos comuns nas consultas: @timestamp, level, service, event.severity, trace_id, message
  • Política ILM: lab-observability-policy (igual à web)
  • Fases do ILM / rollover / idade para exclusão: iguais (hot aos 5 GB/1 d, warm aos 2 d, delete aos 30 d)
  • Pipelines de ingestão: default-enrichment → app-log-enrichment
  • Processadores do pipeline:
    • set event.ingested
    • set event.severity = {{level}} — copiar level
    • lowercase event.severity
    • dissect on message (raramente corresponde; deixa _tmp.*)
    • remove _tmp* — limpeza
  • Campos enriquecidos: event.severity, event.ingested

Projeto proposto da tabela do ClickHouse

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;
  • Chave de partição: toYYYYMM(Timestamp) — mesma justificativa.
  • ORDER BY: (ServiceName, SeverityText, Timestamp) — "mostre os erros do serviço X nos últimos 10 minutos" é a consulta predominante.
  • TraceId / SpanId: Colunas de nível superior com índice de salto bloom_filter em TraceId. NÃO fazem parte da chave primária — isso destruiria a localidade dos dados nas varreduras por intervalo de tempo. O filtro de Bloom acelera buscas pontuais por ID de trace sem afetar a ordenação.
  • Derivação de severidade: Uma coluna MATERIALIZED lowerUTF8(LogAttributes['level']) reproduziria o processador lowercase do ES. Neste laboratório, armazenamos diretamente como SeverityText no coletor (severityparser do OTel).
  • dissect: Ignore-o — raramente corresponde. Se for necessária uma análise estruturada, use o operador regex_parser do OTel, não uma MV do CH (isso mantém o esquema limpo).
  • TTL: 30 dias, igual ao ES.
  • Fases do ILM desnecessárias: igual a web_access — todas, exceto delete.

3. Fluxo de dados: logs-infrastructure-lab

Estado atual

  • Total de documentos: ~18,6 milhões
  • Índices subjacentes: 2
  • Tamanho dos shards primários: ~3,7 GB
  • Shards/réplicas: 2 primários / 1 réplica
  • Campos do mapeamento: ~35 campos de usuário
  • Campos de alta cardinalidade: message (não estruturado), log_message (analisado), pid
  • Campos de baixa cardinalidade: hostname (~10 valores — k8s-node-01..10), process (~10 valores)
  • Campos comuns nas consultas: @timestamp, hostname, process, event.severity, log_message
  • Política ILM / fases / idade para exclusão: iguais (lab-observability-policy)
  • Pipelines de ingestão: default-enrichment → infra-log-parsing
  • Processadores do pipeline:
    • set event.ingested
    • grok %{SYSLOGTIMESTAMP}%{HOSTNAME}%{WORD:process}[...] — extrair 4 campos
    • set event.severity = "info"
    • script (Painless) — elevar a severidade para warn/error com base em palavras-chave de log_message
  • Campos enriquecidos: syslog_timestamp, hostname, process, pid, log_message, event.severity, event.ingested

Projeto proposto da tabela do ClickHouse

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;
  • Partição/ORDER BY: Os dashboards filtram por Hostname → Process → tempo. Ambos têm baixa cardinalidade, portanto compactam muito bem, e a chave de ordenação oferece excelente poda por intervalo de tempo.
  • Substituição do grok: Transfira a regex ao OTel Collector (operador regex_parser) durante a ingestão. Isso mantém a análise próxima à origem e permite ao ClickHouse armazenar campos já estruturados. Alternativa: uma coluna MATERIALIZED com extractAllGroupsVertical(), mas ela paga o custo de CPU em cada inserção no ClickHouse — é melhor pagá-lo uma vez no coletor.
  • Severidade: Uma coluna MATERIALIZED multiIf(positionCaseInsensitive(LogMessage, 'error') > 0, 'error', positionCaseInsensitive(LogMessage, 'warn') > 0, 'warn', 'info') reproduz exatamente o script Painless.
  • TTL: 30 dias.
  • Fases do ILM desnecessárias: mesma situação — apenas delete é necessário.

4. Fluxos de dados: traces-apm-* e logs-apm.* (Carga de trabalho 2 — OTel Demo)

4a. Traces: traces-apm-*

  • Total de documentos: ~13,8 milhões no dia 1, crescendo para ~14 milhões/dia enquanto o gerador de carga é executado (1 índice subjacente por dia)
  • Índices subjacentes: 2–3, dependendo do tempo de execução
  • Campos folha únicos no mapeamento: ~200 (recurso OTel + atributos de span + campos de controle do Elastic APM)
  • Valores distintos de service.name: 19 (16 microsserviços da OTel Demo + frontend-proxy, frontend-web, sample-order-app)
  • Distribuição de processor.event: ~50 % transaction, ~50 % span
  • Linguagens de instrumentação (service.language.name, 10 valores, incluindo unknown para spans emitidos pelo Envoy): nodejs, cpp, python, dotnet, rust, java, php, ruby, go, unknown
  • Cardinalidade de transaction.name: ~7.000 nomes distintos (principalmente combinações de verbo HTTP + rota)
  • Percentis de transaction.duration.us: p50 ≈ 2,5 ms · p95 ≈ 75 ms · p99 ≈ 1,4 s (cauda longa decorrente de chamadas entre serviços e pontos críticos induzidos pelo gerador de carga)
  • event.outcome: ~99,5 % de sucesso · ~0,5 % de falha (as feature flags do flagd usam off por padrão, portanto não há injeção de falhas)

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

  • Fluxos de dados logs-apm.*: 19 no total (18 logs-apm.app.<service>-default, um por serviço, + 1 logs-apm.error-default)
  • Total de documentos em todos os logs-apm.*: dezenas de milhões — cresce rapidamente (só frontend_web chega a ~28 milhões após um dia de carga)
  • 3 principais serviços por volume de logs: frontend_web (~27,8 milhões) · frontend_proxy (~2,0 milhões) · product_catalog (~1,1 milhão). Faixa seguinte: cart (~1 milhão), currency (~400 mil), recommendation (~300 mil). A maioria dos fluxos por serviço permanece abaixo de 100 mil documentos.
  • Observação: frontend_web domina por um fator de 10 porque o Next.js emite um log a cada renderização de página. Considere isso ao dimensionar a tabela otel_logs de destino — o volume por serviço é extremamente desigual.

4c. Destino proposto no ClickHouse: otel_traces

Use o esquema padrão do clickhouseexporter do OTel Collector. Ele faz o mapeamento direto do OTLP, processa lotes e migração de esquema, e não há motivo para escrever a DDL manualmente no caso padrão.

-- 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) — corresponde às consultas típicas "mostre a latência de frontend.checkout nos últimos 15 minutos". É o padrão do exportador.
  • Etapa obrigatória após a criação: o esquema padrão do exportador vem sem índice de salto em TraceId, portanto buscas pontuais por ID de trace (Consulta 5 do Exercício 2B) fariam uma varredura completa. Execute uma vez:
    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 dias, correspondendo ao ILM de origem.

4d. Destino proposto no ClickHouse: otel_logs (logs de aplicativos APM)

Uma tabela otel_logs, não 19.

O Elasticsearch divide por serviço porque cada índice subjacente adiciona sobrecarga de mapeamento por índice, mas oferece armazenamento restrito a cada índice e ILM granular. O ClickHouse inverte essa relação — uma única tabela:

  • Compacta ServiceName como LowCardinality(String) em ~1 byte/linha, independentemente do número de serviços distintos.
  • Permite que um ORDER BY (ServiceName, ...) faça poda gratuita por serviço em todas as consultas.
  • Evita a sobrecarga de partes/mesclagens/processos em segundo plano de 19 tabelas separadas.
  • Possibilita o simples JOIN otel_logs USING (TraceId) otel_traces entre serviços para correlação log↔trace.
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) — atende a "erros do serviço X nos últimos N minutos", a consulta predominante.
  • Correlação de TraceId preservada por uma coluna de nível superior + filtro de Bloom. Logs e traces compartilham tanto o valor de TraceId quanto o nome da coluna, portanto o join ocupa uma única linha.
  • TTL: 30 dias.

5. Linha de base da latência das consultas

Mediana medida em 3 execuções de cada consulta no cluster ES em execução em localhost:9200.

Consulta.took do ES (mediana, ms)Observações
Q1 — 10 principais caminhos de solicitação (status 200)1Primeira execução com cache frio: 200–1.500 ms; execuções com cache aquecido caem para ~1 ms. Agregação de termos em request_path.keyword.
Q2 — Contagem de 5xx por 1 min na última 1 h4A janela de tempo menor + cardinalidade mais baixa tornam a consulta barata mesmo com cache frio.
Q3 — Busca de trace por trace.id1Cache frio: 100–200 ms; aquecido: ~1 ms. Busca pontual em índice invertido — praticamente gratuita quando o segmento do índice está na memória.

(Os valores típicos da primeira execução com cache frio são Q1=254 ms · Q2=6 ms · Q3=158 ms. Execute cada consulta duas vezes antes de registrar, para medir o desempenho estável com cache aquecido.)

Observação didática: esses números variam muito conforme o estado do cache. O importante é o padrão: Q3 é praticamente gratuita no ES porque todo valor de trace.id está em um índice invertido sempre ativo. No ClickHouse, precisamos conquistar esse desempenho com um índice de salto bloom_filter (consulte a Decisão 3 da solução do ADR). Sem ele, WHERE TraceId = ? faria uma varredura completa da tabela otel_traces.

A Parte 3 repetirá essas três consultas no ClickHouse. Espere que Q1 e Q2 sejam mais rápidas no ClickHouse (a varredura colunar de uma pequena coluna Status materializada supera a interseção das listas de postings) e que Q3 fique próxima do ES após a criação do índice de filtro de Bloom. Se Q3 levar mais de 10 vezes o tempo do ES, o filtro de Bloom provavelmente não foi materializado nos grânulos históricos — revise a Seção 4c.


6. Observações gerais

  • Armazenamento total do ES nos 3 fluxos: ~16 GB de primários (web 8,5 + app 3,9 + infra 3,7). O ClickHouse geralmente compacta logs em uma proporção ~10–15× melhor; espere 1–2 GB após a migração para o mesmo conjunto de dados.
  • Processadores que poderiam ficar no lado do cliente: geoip, user_agent, grok, set event.ingested, derivação de severidade — todos eles podem ficar no OTel Collector (processadores: transform, geoip, user_agent, regex_parser, attributes). Fazer o enriquecimento na borda mantém o esquema do backend mínimo. A contrapartida: o custo de CPU dos coletores aumenta com o tamanho da frota, enquanto dictGet() baseado em dicionário do CH é executado centralmente.
  • Veredito sobre as ações do ILM:
    • rollover (tamanho/idade) — Desnecessário. O ClickHouse tem uma tabela com partições; não há índices rotativos para gerenciar.
    • shrink (redução de shards) — Desnecessário. Shards lógicos escalam automaticamente.
    • forcemerge — Desnecessário. As mesclagens em segundo plano são automáticas.
    • set_priority — Desnecessário. Não existe o conceito de prioridade por camada de nós.
    • migrate (roteamento entre camadas de dados hot→warm→cold) — Desnecessário. O ClickHouse Cloud armazena tudo em armazenamento de objetos com cache local automático de leitura; não há camadas de nós entre as quais migrar.
    • delete aos 30 d — Ainda necessário. Reproduzido por TTL.
  • Limitação do armazenamento a 30 % do uso do ES: remova size (numérico de largura fixa, fácil de calcular se necessário), remova referer (baixo valor para consultas, alta cardinalidade), remova o Body de stack traces após 7 dias usando TTL … RECOMPRESS ou TTL por coluna (TTL Timestamp + INTERVAL 7 DAY DELETE WHERE SeverityText = 'info') e reduza a resolução de @timestamp para DateTime (4 bytes) em vez de DateTime64(9) (8 bytes) nos fluxos que não precisam de precisão abaixo de um segundo.

Nesta página

PT