Elasticsearch MigrationClickHouse Workshops

Corrigé de la fiche de modélisation des données

Corrigé type de la fiche de mise en correspondance des modèles de données Elasticsearch et ClickHouse.

Remarque : le nombre de documents, les tailles et les noms d'index différeront dans votre environnement. Les valeurs ci-dessous proviennent d'une exécution de référence après environ 5 jours de génération de données. L'important est la forme des réponses : types, cardinalités, objectif de l'enrichissement et raisonnement qui sous-tend la conception des tables.


1. Flux de données : logs-web_access-lab

État actuel

  • Nombre total de documents : environ 31 millions
  • Nombre d'index sous-jacents : 3 (génération 3, avec un roulement par jour ou par tranche de 5 Go)
  • Taille totale des shards primaires : environ 8,5 Go (sur les 3 index sous-jacents)
  • Shards par index sous-jacent : 2 primaires (issus du modèle de composant lab-logs-settings)
  • Réplicas par index sous-jacent : 1 réplica configuré, à l'état UNASSIGNED dans l'atelier à nœud unique (comportement attendu, puisqu'il n'existe aucun deuxième nœud pour l'héberger)
  • Champs uniques dans le mapping : environ 69 champs utilisateur (hors méta-champs ES comme _id et _index). Le décompte inclut les sous-champs geo.*, user_agent_parsed.* et .keyword ajoutés dynamiquement.
  • Champs à forte cardinalité : remote_addr (environ 65 000), request_path (environ 200 avec une distribution zipfienne, donc modérée), trace.id (si joint à APM), user_agent (plusieurs centaines)
  • Champs utilisés dans la plupart des requêtes (d'après les tableaux de bord de la partie 1) : @timestamp, status, request_path, request_type, service, geo.country_name, user_agent_parsed.name, event.severity, run_time
  • Politique ILM : lab-observability-policy
  • Phases ILM : chaud (priorité 100, roulement à max_age=1d ou max_size=5gb), tiède (à 2 j : réduction à 1 shard, fusion forcée à 1 segment, priorité 50), suppression (à 30 j : suppression + suppression de l'instantané)
  • Condition de roulement : max_size=5 Go OU max_age=1 j
  • Délai avant suppression : 30 jours
  • Pipelines d'ingestion : default-enrichment (par défaut, défini sur tous les flux de données) → web-access-enrichment (appliqué par le routage conditionnel de Filebeat)
  • Processeurs du pipeline :
    • set event.ingested = _ingest.timestamp (default-enrichment) — ajoute l'horodatage d'ingestion
    • geoip on remote_addr → geo.* — ajoute le nom du pays, le nom de la ville et le point géographique
    • user_agent on user_agent → user_agent_parsed.* — analyse le nom, la version, le système d'exploitation et l'appareil
    • set event.severity = "info" — gravité par défaut
    • script (Painless) — remplace la gravité par warn pour les réponses 4xx et par error pour les réponses 5xx
  • Champs enrichis : geo.country_name, geo.city_name, geo.location, user_agent_parsed.name/version/os.*/device.name, event.severity, event.ingested

Conception proposée de la table 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;
  • Clé de partition : toYYYYMM(Timestamp) — une partition par mois. Ce niveau de granularité permet des suppressions efficaces ; un niveau trop fin (toDate) créerait des milliers de parties. Il correspond à la conservation de 30 jours.
  • ORDER BY : (ServiceName, Status, Timestamp) — les requêtes des tableaux de bord (taux de requêtes par service, nombre de réponses 5xx) filtrent d'abord par ServiceName, puis par Status. En dernière position, Timestamp permet d'éliminer des plages temporelles au sein de chaque bloc de tri (service, état).
  • Colonnes de premier niveau ou Map : placez les champs employés par les requêtes fréquentes des tableaux de bord dans des colonnes matérialisées (Status, RequestPath, RunTime) ; conservez les champs rarement interrogés ou de diagnostic (referer, size, user_agent_parsed.device.name) dans LogAttributes. Les colonnes matérialisées sont calculées lors de l'insertion et stockées comme de véritables colonnes, pour bénéficier pleinement des performances de l'analyse colonnaire.
  • Enrichissement GeoIP : dictionnaire (disposition IP_TRIE sur un CSV MaxMind GeoLite2) + dictGet() dans une colonne matérialisée. En effectuant l'enrichissement dans le moteur de stockage, les agents d'ingestion restent sans état et les modifications de schéma ne nécessitent pas le redéploiement des collecteurs.
  • Analyse de l'agent utilisateur : soit (a) le processeur OTel user_agent (solution simple qui émet le nom, le système d'exploitation et l'appareil analysés comme attributs), soit (b) une colonne matérialisée avec extract() / une expression régulière. Le processeur OTel est plus propre pour la migration, car il correspond directement au processeur user_agent existant d'Elastic.
  • TTL : Timestamp + INTERVAL 30 DAY DELETE. Correspond au délai de suppression ILM.
  • Phases ILM inutiles dans CH Cloud : roulement (aucun index successif), réduction (mise à l'échelle automatique des shards logiques), fusion forcée (les fusions en arrière-plan s'en chargent automatiquement), définition de la priorité (aucun cache par niveau n'est exposé) et migration de niveau (les niveaux chaud/tiède/froid ne sont pas pertinents puisque toutes les données résident dans un stockage objet avec mise en cache automatique).

2. Flux de données : logs-application-lab

État actuel

  • Nombre total de documents : environ 12 millions
  • Index sous-jacents : 2
  • Taille des shards primaires : environ 3,9 Go
  • Shards/réplicas : 2 primaires / 1 réplica (réplica non alloué, car il n'existe qu'un nœud)
  • Champs du mapping : environ 42 champs utilisateur
  • Champs à forte cardinalité : trace_id, span_id, message, error.stack (lorsqu'il est présent)
  • Champs fréquemment interrogés : @timestamp, level, service, event.severity, trace_id, message
  • Politique ILM : lab-observability-policy (identique à celle du web)
  • Phases ILM / roulement / délai de suppression : identiques (chaud à 5 Go/1 j, tiède à 2 j, suppression à 30 j)
  • Pipelines d'ingestion : default-enrichment → app-log-enrichment
  • Processeurs du pipeline :
    • set event.ingested
    • set event.severity = {{level}} — copie le niveau
    • lowercase event.severity
    • dissect on message (trouve rarement une correspondance ; laisse des champs _tmp.*)
    • remove _tmp* — nettoyage
  • Champs enrichis : event.severity, event.ingested

Conception proposée de la table 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;
  • Clé de partition : toYYYYMM(Timestamp) — même raisonnement.
  • ORDER BY : (ServiceName, SeverityText, Timestamp) — « afficher les erreurs du service X au cours des 10 dernières minutes » est la requête dominante.
  • TraceId / SpanId : colonnes de premier niveau, assorties d'un index de saut bloom_filter sur TraceId. Elles ne figurent PAS dans la clé primaire, car cela détruirait la localité des données pour les analyses par plage temporelle. Le filtre de Bloom accélère les recherches ponctuelles par identifiant de trace sans influer sur l'ordre de tri.
  • Dérivation de la gravité : une colonne MATERIALIZED lowerUTF8(LogAttributes['level']) reproduirait le processeur ES lowercase. Pour cet atelier, nous la stockons directement dans SeverityText au niveau du collecteur (severityparser OTel).
  • dissect : ignorez-le, car il trouve rarement une correspondance. Si une analyse structurée est nécessaire, utilisez l'opérateur OTel regex_parser, et non une MV CH (le schéma reste ainsi propre).
  • TTL : 30 jours, comme dans ES.
  • Phases ILM inutiles : comme pour web_access, toutes sauf delete.

3. Flux de données : logs-infrastructure-lab

État actuel

  • Nombre total de documents : environ 18,6 millions
  • Index sous-jacents : 2
  • Taille des shards primaires : environ 3,7 Go
  • Shards/réplicas : 2 primaires / 1 réplica
  • Champs du mapping : environ 35 champs utilisateur
  • Champs à forte cardinalité : message (non structuré), log_message (analysé), pid
  • Champs à faible cardinalité : hostname (environ 10 valeurs, k8s-node-01..10), process (environ 10 valeurs)
  • Champs fréquemment interrogés : @timestamp, hostname, process, event.severity, log_message
  • Politique / phases ILM / délai de suppression : identiques (lab-observability-policy)
  • Pipelines d'ingestion : default-enrichment → infra-log-parsing
  • Processeurs du pipeline :
    • set event.ingested
    • grok %{SYSLOGTIMESTAMP}%{HOSTNAME}%{WORD:process}[...] — extrait 4 champs
    • set event.severity = "info"
    • script (Painless) — fait passer la gravité à warn/error selon les mots-clés de log_message
  • Champs enrichis : syslog_timestamp, hostname, process, pid, log_message, event.severity, event.ingested

Conception proposée de la table 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;
  • Partition/ORDER BY : les tableaux de bord filtrent par Hostname → Process → temps. Les deux premières colonnes présentent une faible cardinalité ; elles se compressent donc extrêmement bien, et la clé de tri permet une excellente élimination des plages temporelles.
  • Remplacement de grok : placez l'expression régulière dans l'OTel Collector (opérateur regex_parser) lors de l'ingestion. L'analyse reste proche de la source et ClickHouse stocke des champs déjà structurés. Autre possibilité : une colonne MATERIALIZED avec extractAllGroupsVertical(), mais le processeur de ClickHouse doit alors supporter le coût à chaque insertion ; mieux vaut ne le payer qu'une fois dans le collecteur.
  • Gravité : une colonne MATERIALIZED multiIf(positionCaseInsensitive(LogMessage, 'error') > 0, 'error', positionCaseInsensitive(LogMessage, 'warn') > 0, 'warn', 'info') reproduit exactement le script Painless.
  • TTL : 30 jours.
  • Phases ILM inutiles : même constat ; seule delete est nécessaire.

4. Flux de données : traces-apm-* et logs-apm.* (charge de travail 2 — OTel Demo)

4a. Traces : traces-apm-*

  • Nombre total de documents : environ 13,8 millions le premier jour, puis environ 14 millions par jour tant que le générateur de charge fonctionne (1 index sous-jacent par jour)
  • Index sous-jacents : 2 à 3 selon la durée d'exécution
  • Champs feuilles uniques dans le mapping : environ 200 (attributs de ressource et de span OTel + champs de gestion d'Elastic APM)
  • Valeurs distinctes de service.name : 19 (16 microservices OTel Demo + frontend-proxy, frontend-web, sample-order-app)
  • Répartition de processor.event : environ 50 % de transaction, environ 50 % de span
  • Langages d'instrumentation (service.language.name, 10 valeurs dont unknown pour les spans émis par Envoy) : nodejs, cpp, python, dotnet, rust, java, php, ruby, go, unknown
  • Cardinalité de transaction.name : environ 7 000 noms distincts (principalement des combinaisons verbe HTTP + route)
  • Centiles de transaction.duration.us : p50 ≈ 2,5 ms · p95 ≈ 75 ms · p99 ≈ 1,4 s (longue traîne due aux appels entre services et aux points chauds provoqués par le générateur de charge)
  • event.outcome : environ 99,5 % de succès · environ 0,5 % d'échecs (les indicateurs de fonctionnalité flagd sont réglés sur off par défaut, donc aucune injection de panne)

4b. Journaux : logs-apm.app.* et logs-apm.error-default

  • Flux de données logs-apm.* : 19 au total (18 flux logs-apm.app.<service>-default, un par service, + 1 flux logs-apm.error-default)
  • Nombre total de documents dans tous les flux logs-apm.* : plusieurs dizaines de millions, en croissance rapide (frontend_web atteint à lui seul environ 28 millions après une journée de charge)
  • 3 premiers services par volume de journaux : frontend_web (environ 27,8 millions) · frontend_proxy (environ 2 millions) · product_catalog (environ 1,1 million). Niveau suivant : cart (environ 1 million), currency (environ 400 000), recommendation (environ 300 000). La plupart des flux propres à un service restent sous les 100 000 documents.
  • Observation : frontend_web domine avec un facteur 10, car Next.js émet un journal à chaque rendu de page. Tenez-en compte lors du dimensionnement de la table otel_logs cible : le volume est extrêmement déséquilibré entre les services.

4c. Cible ClickHouse proposée : otel_traces

Utilisez le schéma par défaut du clickhouseexporter de l'OTel Collector. Il correspond directement à OTLP, prend en charge la mise en lots et la migration du schéma ; il n'est donc pas utile d'écrire le DDL à la main dans le cas courant.

-- 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) — correspond aux requêtes courantes telles que « afficher la latence de frontend.checkout au cours des 15 dernières minutes ». Valeur par défaut de l'exportateur.
  • Étape obligatoire après la création : le schéma par défaut de l'exportateur est fourni sans index de saut sur TraceId ; les recherches ponctuelles par identifiant de trace (requête 5 de l'exercice 2B) analyseraient donc toute la table. Exécutez une fois :
    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 jours, comme la politique ILM source.

4d. Cible ClickHouse proposée : otel_logs (journaux d'application APM)

Une seule table otel_logs, et non 19.

Elasticsearch sépare les données par service, car chaque index sous-jacent ajoute un coût de mapping propre à l'index, tout en offrant un stockage précisément délimité par index et une politique ILM fine. ClickHouse inverse ce compromis : une seule table :

  • compresse ServiceName comme LowCardinality(String) à environ 1 octet par ligne, quel que soit le nombre de services distincts ;
  • permet à un unique ORDER BY (ServiceName, ...) d'éliminer gratuitement les données des autres services pour chaque requête ;
  • évite les coûts de parties, de fusion et de tâches en arrière-plan propres à 19 tables distinctes ;
  • autorise une simple JOIN otel_logs USING (TraceId) otel_traces entre les services pour corréler journaux et traces.
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) — couvre la requête dominante : « erreurs du service X au cours des N dernières minutes ».
  • Corrélation par TraceId préservée au moyen d'une colonne de premier niveau et d'un filtre de Bloom. Les journaux et les traces partagent la valeur et le nom de colonne TraceId ; la jointure tient donc sur une ligne.
  • TTL : 30 jours.

5. Mesure de référence de la latence des requêtes

Médiane de 3 exécutions de chaque requête sur le cluster ES actif à l'adresse localhost:9200.

Requête.took d'ES (médiane, ms)Remarques
Q1 — 10 principaux chemins de requête (état 200)1Première exécution à cache froid : 200 à 1 500 ms ; les exécutions à chaud descendent à environ 1 ms. Agrégation terms sur request_path.keyword.
Q2 — Nombre de réponses 5xx par minute au cours de la dernière heure4La fenêtre temporelle plus courte et la cardinalité plus faible rendent cette requête peu coûteuse, même à froid.
Q3 — Recherche de trace par trace.id1À froid : 100 à 200 ms ; à chaud : environ 1 ms. Recherche ponctuelle dans l'index inversé, presque gratuite lorsque le segment d'index est en mémoire.

(Les valeurs de la première exécution avec un cache froid ressemblent généralement à Q1=254 ms · Q2=6 ms · Q3=158 ms. Exécutez chaque requête deux fois avant la mesure afin d'évaluer les performances stabilisées à chaud.)

Note pédagogique : ces valeurs varient énormément selon l'état du cache. L'important est leur forme : Q3 est presque gratuite dans ES, car chaque valeur de trace.id réside dans un index inversé toujours actif. Dans ClickHouse, nous devons obtenir cette performance au moyen d'un index de saut bloom_filter (voir la décision 3 du corrigé de l'ADR). Sans lui, une clause WHERE TraceId = ? analyserait toute la table otel_traces.

La partie 3 rejouera ces trois requêtes dans ClickHouse. Q1 et Q2 devraient être plus rapides dans ClickHouse (l'analyse colonnaire d'une petite colonne matérialisée Status est plus rapide que l'intersection de publications) ; Q3 devrait être proche d'ES une fois l'index avec filtre de Bloom en place. Si Q3 prend plus de 10 fois le temps d'ES, le filtre de Bloom n'a probablement pas été matérialisé sur les granules historiques : revérifiez la section 4c.


6. Observations transversales

  • Stockage ES total des 3 flux : environ 16 Go sur les shards primaires (web 8,5 + application 3,9 + infrastructure 3,7). ClickHouse compresse généralement les journaux 10 à 15 fois mieux ; comptez 1 à 2 Go après la migration pour le même jeu de données.
  • Processeurs qui pourraient s'exécuter côté client : geoip, user_agent, grok, set event.ingested, dérivation de la gravité ; tous peuvent s'exécuter dans l'OTel Collector (processeurs : transform, geoip, user_agent, regex_parser, attributes). L'enrichissement en périphérie permet de réduire le schéma du système cible. En contrepartie, le coût processeur des collecteurs augmente avec la taille du parc, tandis qu'un dictGet() fondé sur un dictionnaire CH s'exécute de façon centralisée.
  • Verdict concernant les actions ILM :
    • rollover (taille/âge) — inutile. ClickHouse utilise une seule table comportant des partitions ; aucun index successif à gérer.
    • shrink (réduction des shards) — inutile. Les shards logiques se mettent à l'échelle automatiquement.
    • forcemerge — inutile. Les fusions en arrière-plan sont automatiques.
    • set_priority — inutile. Il n'existe aucune notion de priorité entre niveaux de nœuds.
    • migrate (routage des données chaud→tiède→froid) — inutile. ClickHouse Cloud stocke toutes les données dans un stockage objet assorti d'un cache local automatique en lecture ; il n'existe aucun niveau de nœuds entre lesquels les migrer.
    • delete à 30 j — toujours nécessaire. Reproduit par TTL.
  • Plafonner le stockage à 30 % de l'utilisation d'ES : supprimez size (valeur numérique de largeur fixe, facile à recalculer au besoin), supprimez referer (peu utile pour les requêtes, forte cardinalité), supprimez le Body des traces de pile après 7 jours au moyen d'un TTL … RECOMPRESS ou d'un TTL par colonne (TTL Timestamp + INTERVAL 7 DAY DELETE WHERE SeverityText = 'info') et réduisez la précision d'@timestamp à DateTime (4 octets) au lieu de DateTime64(9) (8 octets) pour les flux qui ne nécessitent pas une précision inférieure à la seconde.

Sur cette page

FR