Elasticsearch MigrationClickHouse Workshops

Soluciones de SQL avanzado

Respuestas modelo para seis consultas de ClickHouse que sustituyen flujos difíciles de Elasticsearch.

Base de datos: Todas las tablas están en otel. Ejecuta clickhouse client --database otel ... o primero USE otel;.


Ejercicio 1: JOIN entre señales — Logs + Trazas

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;

Puntos didácticos:

  • Esta consulta sustituye un flujo de tres pasos en Kibana: buscar las trazas más lentas, copiar sus ID y buscar logs para cada una.
  • ClickHouse ejecuta la CTE una vez y usa un hash join para comparar l.TraceId con las 10 filas.
  • El índice de salto bloom_filter en TraceId de otel_logs_v2 acelera cada sondeo.
  • No se usa groupArray() porque queremos filas individuales; el Ejercicio 6 lo usa para salida agregada. Allí, groupArray() conserva muestras por grupo.

¿Por qué las dos cláusulas WHERE adicionales? Sin SpanKind = 'Server', las 10 trazas más lentas son long-polls gRPC internos EventStream de flagd, abiertos ~10 minutos (Duration ≈ 600,000,000,000 ns) y sin logs asociados; el JOIN no devuelve nada. Limitar a solicitudes entrantes y excluir flagd muestra solicitudes reales como frontend-proxy / ingress de ~8–9 segundos con logs relacionados. «Top N por latencia» casi siempre necesita un filtro por tipo/nombre.

Valores del enum: El exportador escribe SpanKind como 'Server', 'Client', 'Internal', 'Producer', 'Consumer', no nombres proto como 'SPAN_KIND_SERVER'. Compruébalo con SELECT DISTINCT SpanKind FROM otel_traces. El filtro operativo que se conserva es SpanKind = 'Server'.

Forma esperada: Varias filas por traza, ordenadas por duración y marca temporal.


Ejercicio 2: funciones de ventana — detección de anomalías

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;

Puntos didácticos:

  • LAG(error_rate_pct) OVER (PARTITION BY ServiceName ORDER BY minute) compara cada minuto con el anterior del mismo servicio. Sin PARTITION BY, mezclaría servicios.
  • En Elasticsearch harían falta date_histogram, JSON en el cliente y deltas en Python/JavaScript: al menos 2 llamadas, código y estado del bucket anterior.
  • La CTE with minute_errors AS (...) calcula métricas y la consulta externa aplica la ventana; ClickHouse encadena ambas sin materializar todo el resultado.
  • HAVING total > 50 evita que 1 error entre 3 solicitudes (33 %) domine. Filtra siempre por tamaño de muestra antes de comparar tasas.

¿Por qué minuto / 1,2× / 1 hora? Los generadores mantienen ~5 % de errores con ruido Poisson. Empíricamente:

UmbralPares por minuto que lo superan (última 1 h, ~305 muestras)
pico > 1.10 (10%)40
pico > 1.20 (20%)4
pico > 1.50 (50%)0
pico > 2.00 (2×)0
pico > 3.00 (3×)0

Un salto del 20 % muestra sucesos raros sin el ruido habitual:

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

En producción se ajustaría a 1,5×–3× sobre 5 minutos. Una regla horaria de 3× es demasiado gruesa para estos datos sintéticos, pero puede ser adecuada con tráfico humano real.


Ejercicio 3: GROUP BY ilimitado — inventario 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;

Puntos didácticos:

  • Sin LIMIT, ClickHouse devuelve cada RequestPage: 100, 10 000 o 1 000 000. Terms de Elasticsearch con size solo devuelve N.
  • La agregación composite de ES pagina con after_key; cada página es otra llamada. ClickHouse hace una pasada.
  • quantile(0.95)(toFloat64OrZero(...)) calcula p95 dentro del GROUP BY, sin subagregación aparte.
  • toFloat64OrZero() maneja valores vacíos o no numéricos devolviendo 0.
  • La clave original del Map sigue siendo run_time.

Ejercicio 4: detección de secuencias — flujos que acaban en error

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;

Puntos didácticos:

  • (?1)(?t<=10).*(?2)(?t<=10).*(?3) es un patrón de secuencia:
    • (?N) coincide con la condición N
    • (?t<=10) limita la diferencia hasta el siguiente evento a ≤ 10 segundos
    • .* admite eventos intermedios
    • Las condiciones 1, 2 y 3 deben ocurrir en orden, cada una en 10 segundos
  • Opera en filas agrupadas por RemoteAddr como secuencia ordenada.
  • sequenceMatch() es propio de ClickHouse; ni DSL ni ES|QL expresan patrones ordenados entre servicios.
  • (?t op N) admite <, <=, ==, >=, >; windowFunnel() recibe la ventana aparte e indica hasta dónde llegó cada fila.

Salida verificada: Sin (?t<=10) ('(?1).*(?2).*(?3)'), coinciden ~50 de 51 IP porque todas acaban siguiendo la ruta en una hora. El límite reduce a ~9 IP y produce una señal más útil.

Si devuelve 0 filas: Ejecuta SELECT DISTINCT ServiceName FROM otel_logs_v2 WHERE RequestType != ''; el laboratorio usa api-gateway, order-service, payment-service, inventory-service, web-frontend.


Ejercicio 5: agregación condicional — salud del servicio

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;

Puntos didácticos:

  • El combinador -If se añade a cualquier agregado: countIf, avgIf, sumIf, quantileIf, uniqIf, minIf, maxIf, etc.
  • La consulta calcula 12 métricas con 7 condiciones en una exploración. En ES cada métrica necesita filter → metric, unas 150 líneas JSON.
  • quantileIf(0.95)(latency, StatusCode < 500) calcula p95 solo para solicitudes correctas, difícil sin bucket_script en ES.
  • uniqIf(RemoteAddr, StatusCode >= 500) cuenta clientes afectados. ES exigiría cardinality filtrado; ambos son aproximados, pero ClickHouse es más simple.
  • El combinador uniqIf mantiene la condición dentro del agregado.

Ejercicio 6: investigación de causa raíz con 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;

Puntos didácticos:

  • Sustituye 3 llamadas ES + unión JSON en cliente: tasa de errores, términos de mensajes e ID de traza.
  • groupArray(N)(expr) recoge hasta N valores de expr por grupo. top_hits de ES es lo más cercano, pero no funciona entre JOIN.
  • Las CTE se calculan una vez y se reutilizan. error_services alimenta tres CTE y el SELECT final.
  • La relación de servicios de interés sigue centralizada en error_services.
  • LEFT JOIN conserva un servicio aunque no tenga ID de traza; INNER JOIN lo eliminaría.
  • HAVING errors > 10 evita ruido de servicios de poco tráfico; ajusta el umbral.

¿Por qué sample_trace_ids está vacío? Los resultados (web-frontend, payment-service, order-service, etc.) proceden de generadores basados en archivos y no propagan contexto; sus filas tienen TraceId = ''. affected_traces filtra por TraceId != '', por lo que groupArray(5)(a.TraceId) devuelve ['','','','',''] debido al LEFT JOIN. SELECT ServiceName, countIf(TraceId != '') AS rows_with_trace, count() AS total FROM otel_logs_v2 GROUP BY ServiceName ORDER BY total DESC muestra 0 para servicios de archivos y valores no nulos para OTel Demo (frontend-proxy, product-catalog, cart, …), que no emite 5xx en StatusCode y por eso no aparece en error_services. En producción, el mismo servicio instrumentado emitiría errores y trazas, y sample_trace_ids permitiría abrir la traza. El patrón SQL es el adecuado aunque el laboratorio no muestre el recorrido completo. En otras palabras, en esos servicios se cumple rows_with_trace = 0 aunque otel_logs_v2 contenga muchas filas.

En esta página

ES