Elasticsearch MigrationClickHouse Workshops

Soluções de SQL avançado

Respostas-modelo para as seis consultas do ClickHouse que substituem fluxos de trabalho difíceis do Elasticsearch.

Banco de dados: Todas as tabelas abaixo ficam no banco de dados otel. Execute primeiro com clickhouse client --database otel ... ou USE otel;.


Exercício 1: JOIN entre sinais — logs + traces

WITH slow_traces AS (
    SELECT TraceId, ServiceName, SpanName, Duration
    FROM otel_traces
    WHERE Timestamp >= now() - INTERVAL 1 HOUR
      AND SpanKind = 'Server'             -- user-facing request spans only
      AND SpanName NOT LIKE '%flagd%'     -- exclude long-poll feature-flag streams
    ORDER BY Duration DESC
    LIMIT 10
)
SELECT
    t.TraceId,
    t.ServiceName       AS trace_service,
    t.SpanName,
    t.Duration          AS trace_duration_ns,
    l.Timestamp         AS log_time,
    l.ServiceName       AS log_service,
    l.SeverityText,
    l.Body
FROM slow_traces t
JOIN otel_logs_v2 l ON t.TraceId = l.TraceId
ORDER BY t.Duration DESC, l.Timestamp ASC;

Pontos de ensino:

  • Esta única consulta substitui um fluxo de trabalho de 3 etapas no Kibana: (1) consultar no APM os traces mais lentos, (2) copiar os IDs dos traces e (3) pesquisar os logs de cada ID individualmente.
  • O ClickHouse executa a CTE uma vez e usa um hash join para comparar l.TraceId com o conjunto de 10 linhas resultante — extremamente eficiente.
  • O índice de skipping bloom_filter em TraceId de otel_logs_v2 acelera a busca do join para cada um dos 10 IDs de traces lentos.
  • groupArray() não foi usado aqui porque queremos linhas de log individuais, não arrays. O Exercício 6 usa groupArray() para a saída agregada.

Por que as duas cláusulas WHERE adicionais? Sem SpanKind = 'Server', os 10 spans mais lentos serão todos long polls gRPC internos EventStream de flagd — spans de controle que permanecem abertos por cerca de 10 minutos (Duration ≈ 600,000,000,000 ns) e não têm nenhum registro de log associado, portanto o JOIN não retorna linhas. Restringir a SpanKind = 'Server' (requisições HTTP recebidas) e excluir flagd pelo nome revela requisições lentas reais voltadas ao usuário, como os spans frontend-proxy / ingress, que levam cerca de 8–9 segundos; cada um deles tem logs de aplicação correspondentes, então o JOIN funciona. A instrumentação do OTel Demo é naturalmente ruidosa no extremo superior da distribuição de duração — um lembrete útil de que os "N principais por latência" quase sempre precisam de um filtro de tipo/nome para serem significativos.

Observação sobre valores de enum: O exporter do collector do OTel do ClickHouse grava SpanKind como 'Server', 'Client', 'Internal', 'Producer', 'Consumer' — não como nomes do proto OTel, como 'SPAN_KIND_SERVER'. Sempre confira o que realmente está armazenado: SELECT DISTINCT SpanKind FROM otel_traces.

Formato de saída esperado: várias linhas por trace (todas as entradas de log que compartilham esse TraceId), ordenadas pela duração do trace e depois pelo timestamp do log.


Exercício 2: funções de janela — detecção de anomalias

WITH minute_errors AS (
    SELECT
        ServiceName,
        toStartOfMinute(Timestamp)  AS minute,
        countIf(StatusCode >= 500)  AS errors,
        count()                     AS total,
        if(total > 0, errors / total * 100, 0) AS error_rate_pct
    FROM otel_logs_v2
    WHERE TimestampTime >= now() - INTERVAL 1 HOUR
      AND RequestType != ''
    GROUP BY ServiceName, minute
    HAVING total > 50         -- skip low-sample minutes where the rate is noise
),
with_lag AS (
    SELECT
        ServiceName,
        minute,
        error_rate_pct,
        LAG(error_rate_pct) OVER (PARTITION BY ServiceName ORDER BY minute) AS prev_minute_rate,
        if(prev_minute_rate > 0, error_rate_pct / prev_minute_rate, 0)      AS spike_ratio
    FROM minute_errors
)
SELECT *
FROM with_lag
WHERE spike_ratio > 1.2
ORDER BY spike_ratio DESC;

Pontos de ensino:

  • LAG(error_rate_pct) OVER (PARTITION BY ServiceName ORDER BY minute) compara cada minuto com o minuto anterior do mesmo serviço. Sem PARTITION BY, você compararia incorretamente serviços diferentes e obteria uma razão sem significado entre serviços.
  • No Elasticsearch, seria necessária uma agregação date_histogram para obter contagens por minuto, baixar o JSON no cliente e calcular as diferenças em Python/JavaScript. Isso exige no mínimo 2 chamadas de API + código da aplicação, além de estado para lembrar o bucket anterior.
  • A CTE with minute_errors AS (...) calcula as métricas básicas; a consulta externa aplica a função de janela. Essa estrutura em duas etapas mantém a intenção clara e muitas vezes é mais eficiente do que uma única consulta complexa — o ClickHouse pode encadear as duas etapas sem materializar todo o resultado intermediário.
  • HAVING total > 50 na primeira CTE é essencial: sem ele, um minuto com 3 requisições e 1 erro (taxa de 33%!) superaria todos os outros buckets. Sempre filtre as agregações pelo tamanho da amostra antes de comparar razões.

Por que esses parâmetros específicos (minuto/1,2×/1 hora)? Os geradores de logs do laboratório emitem deliberadamente uma taxa estável de cerca de 5% de erros por serviço (ruído de Poisson em torno da média). Empiricamente:

LimitePares por minuto que o ultrapassaram (última hora, cerca de 305 amostras)
spike > 1.10 (10%)40
spike > 1.20 (20%)4
spike > 1.50 (50%)0
spike > 2.00 (2×)0
spike > 3.00 (3×)0

Um salto de 20% (1,2×) com resolução por minuto é o menor limite que revela eventos realmente raros, em vez do ruído cotidiano — saída de exemplo típica:

ServiceName    minute               rate_pct  prev_minute_rate  spike_ratio
web-frontend   2026-05-09 04:10:00  5.186     4.263             1.216
web-frontend   2026-05-09 03:48:00  5.501     4.549             1.209
api-gateway    2026-05-09 04:24:00  4.714     3.903             1.208
web-frontend   2026-05-09 04:01:00  5.357     4.460             1.201

Em uma regra de produção, você ajustaria aos próprios padrões de tráfego — normalmente 1,5×–3× em uma janela móvel de 5 minutos. Um bucket de uma hora/3× como regra (limite comum de SRE para um "incidente evidente") é grosseiro demais para os dados deste laboratório, pois a variação horária relativa do gerador sintético é de cerca de 3%; em produção, com tráfego real gerado por usuários, 3× por hora é exatamente o limite adequado para "temos um problema real".


Exercício 3: GROUP BY ilimitado — inventário completo de endpoints

SELECT
    RequestPage,
    count()                                              AS total_requests,
    countIf(StatusCode >= 500)                           AS errors,
    round(errors / total_requests * 100, 2)              AS error_rate_pct,
    quantile(0.95)(toFloat64OrZero(LogAttributes['run_time'])) AS p95_latency
FROM otel_logs_v2
WHERE RequestType != ''
GROUP BY RequestPage
ORDER BY total_requests DESC;

Pontos de ensino:

  • Sem LIMIT, o ClickHouse retorna todos os valores exclusivos de RequestPage — 100, 10.000 ou 1.000.000. Uma agregação Terms do Elasticsearch com size retornaria apenas os N principais.
  • No Elasticsearch, a alternativa de agregação composta pagina os resultados usando after_key. Cada página é uma chamada de API separada. O ClickHouse faz isso em uma única passagem.
  • quantile(0.95)(toFloat64OrZero(...)) calcula a latência p95 inline com o GROUP BY — sem necessidade de subagregação separada. No Elasticsearch, adicionar percentis a uma agregação terms dobra o payload da resposta e a complexidade da consulta.
  • toFloat64OrZero() trata com segurança linhas em que run_time está vazio ou não é numérico (retorna 0 em vez de gerar erro).

Exercício 4: detecção de sequências — fluxos de requisições que levam a erros

SELECT
    RemoteAddr,
    count() AS occurrence_count
FROM otel_logs_v2
WHERE TimestampTime >= now() - INTERVAL 1 HOUR
GROUP BY RemoteAddr
HAVING sequenceMatch('(?1)(?t<=10).*(?2)(?t<=10).*(?3)')(
    TimestampTime,                                       -- sequenceMatch requires DateTime, not DateTime64
    ServiceName = 'api-gateway'   AND StatusCode = 200,
    ServiceName = 'order-service' AND StatusCode = 200,
    ServiceName = 'payment-service' AND StatusCode >= 500
)
ORDER BY occurrence_count DESC
LIMIT 20;

Pontos de ensino:

  • (?1)(?t<=10).*(?2)(?t<=10).*(?3) é um padrão semelhante a uma expressão regular sobre sequências de eventos:
    • (?N) corresponde a um evento que satisfaz a condição N
    • (?t<=10) restringe a diferença de tempo até o próximo evento a ≤ 10 segundos (inclusive)
    • .* corresponde a zero ou mais eventos intermediários entre as duas condições-âncora
    • O padrão completo exige que as condições 1, 2 e 3 ocorram em ordem, e cada condição correspondente seguinte deve chegar em até 10 segundos após a anterior
  • A função opera sobre linhas agrupadas por RemoteAddr, tratando cada grupo como uma sequência ordenada de eventos.
  • sequenceMatch() é exclusivo do ClickHouse. Nem a DSL do Elasticsearch nem ES|QL conseguem expressar padrões de eventos ordenados em várias etapas entre serviços.
  • O operador de diferença temporal (?t op N) aceita <, <=, ==, >=, >. Ele difere de windowFunnel(), que recebe a janela como argumento numérico separado e informa até onde cada linha avançou no funil.

Saída verificada: Sem o limite de tempo (?t<=10) (isto é, '(?1).*(?2).*(?3)'), esta consulta retorna cerca de 50 dos 51 IPs distintos do laboratório — todos os IPs acabam percorrendo o caminho api-gateway → order-service → payment-service-5xx em uma hora; portanto, o padrão corresponde em quase todos os casos e a resposta não discrimina muito. Adicionar o limite de 10 segundos entre eventos correspondentes sucessivos reduz o resultado para cerca de 9 IPs, muito mais próximo do sinal de "anomalia real" desejado em produção.

Se a consulta retornar 0 linhas: os nomes de serviço do gerador de logs podem ser diferentes. Execute SELECT DISTINCT ServiceName FROM otel_logs_v2 WHERE RequestType != '' para encontrar os nomes reais dos serviços de acesso web — o laboratório usa api-gateway, order-service, payment-service, inventory-service e web-frontend.


Exercício 5: agregação condicional — integridade do serviço com várias métricas

SELECT
    ServiceName,
    count()                                                            AS total_events,
    countIf(StatusCode >= 200 AND StatusCode < 300)                   AS success_2xx,
    countIf(StatusCode >= 400 AND StatusCode < 500)                   AS client_errors_4xx,
    countIf(StatusCode >= 500)                                         AS server_errors_5xx,
    round(server_errors_5xx / total_events * 100, 2)                  AS error_rate_pct,
    quantileIf(0.50)(toFloat64OrZero(LogAttributes['run_time']), StatusCode < 500) AS p50_latency_ok,
    quantileIf(0.95)(toFloat64OrZero(LogAttributes['run_time']), StatusCode < 500) AS p95_latency_ok,
    quantileIf(0.50)(toFloat64OrZero(LogAttributes['run_time']), StatusCode >= 500) AS p50_latency_err,
    uniqIf(RemoteAddr, StatusCode >= 500)                             AS unique_affected_ips,
    minIf(Timestamp, StatusCode >= 500)                               AS first_error_at,
    maxIf(Timestamp, StatusCode >= 500)                               AS last_error_at
FROM otel_logs_v2
WHERE TimestampTime >= now() - INTERVAL 1 HOUR
  AND RequestType != ''
GROUP BY ServiceName
ORDER BY error_rate_pct DESC;

Pontos de ensino:

  • O combinador -If pode ser acrescentado a qualquer função de agregação do ClickHouse: countIf, avgIf, sumIf, quantileIf, uniqIf, minIf, maxIf etc.
  • Esta única consulta calcula 12 métricas em 7 condições diferentes com um scan da tabela. No Elasticsearch, cada métrica condicional exige seu próprio aninhamento filter → metric, produzindo cerca de 150 linhas de JSON aninhado para uma consulta equivalente.
  • quantileIf(0.95)(latency, StatusCode < 500) calcula p95 somente para requisições bem-sucedidas — impossível de expressar em uma única agregação do ES sem bucket_script.
  • uniqIf(RemoteAddr, StatusCode >= 500) conta clientes distintos afetados — útil para medir o raio de impacto. No ES, isso exigiria uma subagregação cardinality filtrada com sua aproximação HyperLogLog; uniqIf do ClickHouse também é aproximado (HLL), mas a sintaxe é muito mais simples.

Exercício 6: investigação da causa raiz com CTE

WITH
error_services AS (
    SELECT
        ServiceName,
        countIf(StatusCode >= 500) AS errors,
        count()                    AS total
    FROM otel_logs_v2
    WHERE TimestampTime >= now() - INTERVAL 1 HOUR
      AND RequestType != ''
    GROUP BY ServiceName
    HAVING errors > 10
    ORDER BY errors / total DESC
    LIMIT 3
),
top_errors AS (
    SELECT
        l.ServiceName,
        l.Body,
        count() AS occurrences
    FROM otel_logs_v2 l
    INNER JOIN error_services e ON l.ServiceName = e.ServiceName
    WHERE l.StatusCode >= 500
      AND l.TimestampTime >= now() - INTERVAL 1 HOUR
    GROUP BY l.ServiceName, l.Body
    ORDER BY occurrences DESC
    LIMIT 10
),
affected_traces AS (
    SELECT DISTINCT
        l.TraceId,
        l.ServiceName
    FROM otel_logs_v2 l
    INNER JOIN error_services e ON l.ServiceName = e.ServiceName
    WHERE l.StatusCode >= 500
      AND l.TraceId != ''
      AND l.TimestampTime >= now() - INTERVAL 1 HOUR
    LIMIT 50
)
SELECT
    e.ServiceName,
    e.errors,
    e.total,
    round(e.errors / e.total * 100, 2)  AS error_rate_pct,
    groupArray(10)(t.Body)              AS sample_error_messages,
    groupArray(5)(a.TraceId)            AS sample_trace_ids
FROM error_services e
LEFT JOIN top_errors t       ON e.ServiceName = t.ServiceName
LEFT JOIN affected_traces a  ON e.ServiceName = a.ServiceName
GROUP BY e.ServiceName, e.errors, e.total
ORDER BY error_rate_pct DESC;

Pontos de ensino:

  • Esta única consulta substitui 3 chamadas separadas à API do Elasticsearch + combinação de JSON no cliente: (1) taxa de erros em histograma de datas, (2) agregação terms das mensagens de erro e (3) agregação terms dos IDs de traces.
  • groupArray(N)(expr) coleta até N valores de expr em um array por grupo. É o equivalente do ClickHouse à coleta de valores de amostra — não há equivalente simples no ES (top_hits é o mais próximo, mas funciona somente dentro de subagregações terms, não entre JOINs).
  • As CTEs no ClickHouse são calculadas uma vez e reutilizadas. error_services é referenciada por três CTEs posteriores e pelo SELECT final — o ClickHouse a materializa uma vez.
  • O LEFT JOIN no SELECT final garante que uma linha de serviço apareça mesmo sem IDs de traces correspondentes (por exemplo, se todos os erros não tiverem TraceId). Um INNER JOIN eliminaria silenciosamente esses serviços do resultado.
  • O filtro HAVING errors > 10 em error_services evita ruído de serviços com tráfego muito baixo que por acaso tenham um único erro. Ajuste o limite conforme sua taxa de ingestão.

Por que sample_trace_ids está vazio neste laboratório? A consulta retorna linhas para web-frontend, payment-service, order-service etc. — são os geradores de logs baseados em arquivos apresentados na Parte 1, que não propagam o contexto de trace; assim, todas as linhas desses serviços em otel_logs_v2 têm TraceId = ''. A CTE affected_traces filtra por TraceId != '' e fica vazia, portanto groupArray(5)(a.TraceId) retorna ['','','','',''] (o LEFT JOIN mantém a linha pai, mas não há valores reais de TraceId). Para verificar, SELECT ServiceName, countIf(TraceId != '') AS rows_with_trace, count() AS total FROM otel_logs_v2 GROUP BY ServiceName ORDER BY total DESC mostra que os serviços baseados em arquivos têm rows_with_trace = 0, enquanto os serviços do OTel Demo (frontend-proxy, product-catalog, cart, …) têm valores diferentes de zero — mas estes não emitem atualmente códigos HTTP 5xx pela coluna StatusCode, portanto não aparecem em error_services. Em um ambiente de produção real, a aplicação que emite erros 5xx e a que emite traces seriam o mesmo serviço instrumentado com OTel, e sample_trace_ids seria um ponto de partida clicável para aprofundar a análise do trace. O laboratório não consegue demonstrar isso de ponta a ponta com as fontes de dados atuais, mas o padrão SQL é exatamente o que você executaria em produção.

Nesta página

PT