Déduplication des logs avec Elasticsearch

Les événements en double provenant de services d'application en mauvaise santé rendent la recherche de logs complexe. Découvrez comment gérer les doublons à l'aide de Logstash, Beats et Elastic Agent.

Déduplication des logs avec Elasticsearch

Les SRE sont quotidiennement submergés par de grands volumes de logs provenant d’applications bruyantes. Dans son ouvrage séminal The Mythical Man Month, Frederick P. Brooks a déclaré que « tous les programmeurs sont des optimistes ». Cet optimisme se manifeste par le fait que les ingénieurs logiciels ne mettent pas en place de contrôles pour empêcher leurs applications d’envoyer des logs en continu dans des situations de défaillance exceptionnelles. Dans les grandes organisations dotées de plateformes de logging centralisées, ce flux d’événements est ingéré dans les plateformes de logging, occupant des volumes de stockage et des ressources de calcul considérables. Côté humain, les SRE se sentent dépassés et souffrent de fatigue liée aux alertes, submergés par la vague de messages, un peu comme ceci :

 

![Surfer Engulfed by Waves Gif](./images/1.gif)

 

Les développeurs qui conçoivent des logiciels, notamment des microservices et des applications clés, sont responsables de s’assurer qu’ils n’envoient pas d’événements de log en double et qu’ils envoient les bons événements de log au niveau approprié. Néanmoins, des situations telles que l’utilisation de solutions tierces ou la maintenance de services vieillissants signifient que nous ne pouvons pas toujours garantir l’application de pratiques de logging responsables. Même si nous supprimons les champs inutiles comme expliqué dans cet article sur l’élagage des champs des événements de log entrants, nous avons toujours un problème avec le fait de stocker un grand nombre d’événements en double. Nous abordons ici les défis liés à l’identification des logs en double provenant de services problématiques et la manière de dédupliquer les données à l’aide d’Elastic Beats, de Logstash et d’Elastic Agent.

What is a duplicate log entry?

Before diving into the various ways of preventing these duplicates from making it into your logging platform, we need to understand what a duplicate is. In my prior life as a software engineer, I was responsible for developing and maintaining an ecosystem of vast microservices. Some had considered retry logic that, after some time, would shut the service down gracefully and trigger appropriate alerts. However, not all services are built to gracefully handle these cases. Service misconfiguration can also contribute to event duplication. Inadvertently changing the production log level from WARN to TRACE can lead to more aggressive event volumes that have to be handled by the logging platform.

Elasticsearch automatically generates a unique ID for each document ingested unless a document contains an _id field on ingestion. Therefore, if your service is sending repeated alerts, you run the risk of having the same event stored as multiple documents with different IDs. Another cause can be due to retry mechanisms for the tools used for log collection. A notable example is for Filebeat, where a lost connection or shutdown can cause the retry mechanism of Filebeat to resend an event until the output acknowledges receipt of the event.

Aperçu des outils

Dans ce blog, nous examinerons les outils disponibles dans quatre outils Elastic :

  1. Logstash est un outil de pipeline ETL gratuit et open source qui vous permet l’ingestion, la transformation et la sortie de données entre une myriade de sources, y compris l’ingestion dans et la sortie depuis Elasticsearch. Ces exemples seront utilisés pour montrer le comportement différent d’Elasticsearch lors de la génération de logs avec et sans ID particulier.
  2. Beats sont une famille d’agents de transfert légers qui nous permettent l’ingestion d’événements depuis une source donnée non seulement dans Elasticsearch, mais aussi dans d’autres sorties, notamment Kafka, Redis ou Logstash.
  3. Pipelines d’ingestion permettent d’appliquer des transformations et des enrichissements aux documents ingérés dans Elasticsearch. C’est comme exécuter la partie filter de Logstash directement dans Elasticsearch sans avoir besoin d’exécuter un autre service. De nouveaux pipelines peuvent être créés soit depuis l'écran Stack Management > Ingestion Pipelines, soit via l'API _ingestion, comme indiqué dans la documentation.
  4. Elastic Agent est un agent unique capable de s'exécuter sur votre hôte et d'envoyer des logs, des indicateurs et des données de sécurité provenant de multiples services et de votre infrastructure vers Elasticsearch en utilisant les diverses intégrations prises en charge. Quelle que soit la raison des doublons, il existe plusieurs pistes d'action possibles dans l'écosystème Elastic.

Ingestion sans ID spécifiés

L’approche par défaut consiste à ignorer et à ingérer tous les événements. Lorsqu’aucun ID n’est spécifié sur un document, Elasticsearch génère automatiquement un nouvel ID pour chaque document qu’il reçoit. Prenons un exemple simple, disponible dans ce référentiel GitHub, utilisant un serveur HTTP Express simple. Le serveur, lors de son exécution, expose un point de terminaison unique renvoyant un seul message de log :

{"event":{"transaction_id":1,"data_set":"my-logging-app"},"message":"WARN: Unable to get an interesting response"}

En utilisant Logstash, nous pouvons interroger le point de terminaison http://locahost:3000/ toutes les 60 secondes et envoyer le résultat à Elasticsearch. Notre logstash.conf ressemble à ce qui suit :

input {
  http_poller {
    urls => {
      simple_server => "http://localhost:3000"
    }
    request_timeout => 60
    schedule => { cron => "* * * * * UTC"}
    codec => "json"
  }
}
output {
  elasticsearch { 
    cloud_id => "${ELASTIC_CLOUD_ID}" 
    cloud_auth => "${ELASTIC_CLOUD_AUTH}"
    index => "my-logstash-index"
    }
}

Logstash enverra chaque événement, et en l'absence d'identifiant sur l'événement, Elasticsearch générera un nouvel _id champ servant d'identifiant unique pour chaque document :

GET my-logstash-index/_search
{
  "took": 0,
  "timed_out": false,
  "_shards": {
    "total": 1,
    "successful": 1,
    "skipped": 0,
    "failed": 0
  },
  "hits": {
    "total": {
      "value": 11,
      "relation": "eq"
    },
    "max_score": 1,
    "hits": [
      {
        "_index": "my-logstash-index",
        "_id": "-j83XYsBOwNNS8Sc0Bja",
        "_score": 1,
        "_source": {
          "@version": "1",
          "event": {
            "transaction_id": 1,
            "original": """{"event":{"transaction_id":1,"data_set":"my-logging-app"},"message":"WARN: Unable to get an interesting response"}""",
            "data_set": "my-logging-app"
          },
          "message": "WARN: Unable to get an interesting response",
          "@timestamp": "2023-10-23T15:47:00.528205Z"
        }
      },
      {
        "_index": "my-logstash-index",
        "_id": "NT84XYsBOwNNS8ScuRlO",
        "_score": 1,
        "_source": {
          "@version": "1",
          "event": {
            "transaction_id": 1,
            "original": """{"event":{"transaction_id":1,"data_set":"my-logging-app"},"message":"WARN: Unable to get an interesting response"}""",
            "data_set": "my-logging-app"
          },
          "message": "WARN: Unable to get an interesting response",
          "@timestamp": "2023-10-23T15:48:00.314262Z"
        }
      },
      // Other documents omitted
    ]
  }
}

Ce comportement est cohérent pour Beats, les pipelines d'ingestion et Elastic Agent, car ils enverront tous les événements reçus sans configuration supplémentaire.

BYO-ID

La spécification d'un ID unique pour chaque événement à l'aide d'un ID existant permet de contourner l'étape de génération d'ID par Elasticsearch abordée dans la section précédente. L'ingestion d'un document pour lequel cet attribut existe déjà amènera Elasticsearch à vérifier si un document possédant cet ID existe dans l'index, et à mettre à jour le document le cas échéant. Cela entraîne une surcharge, car l'index doit être recherché pour vérifier si un document possédant le même _id existe. Pour prolonger notre exemple Logstash ci-dessus, il est possible de spécifier la valeur de l'ID de document dans Logstash en définissant l'option document_id dans le plug-in de sortie Elasticsearch, qui serait utilisé pour l'ingestion d'événements dans Elasticsearch :

# http_poller configuration omitted
output {
  elasticsearch { 
    cloud_id => "${ELASTIC_CLOUD_ID}" 
    cloud_auth => "${ELASTIC_CLOUD_AUTH}"
    index => "my-unique-logstash-index"
    document_id => "%{[event][transaction_id]}"
    }
}

Cela définira la valeur du champ _id sur la valeur de event.transaction_id. Dans notre cas, cela signifie que le nouveau document remplacera le document existant lors de l'ingestion, car les deux documents ont un _id de 1 :

GET my-unique-logstash-index/_search
{
  "took": 48,
  "timed_out": false,
  "_shards": {
    "total": 1,
    "successful": 1,
    "skipped": 0,
    "failed": 0
  },
  "hits": {
    "total": {
      "value": 1,
      "relation": "eq"
    },
    "max_score": 1,
    "hits": [
      {
        "_index": "my-unique-logstash-index",
        "_id": "1",
        "_score": 1,
        "_source": {
          "@timestamp": "2023-10-23T16:33:00.358585Z",
          "message": "WARN: Unable to get an interesting response",
          "@version": "1",
          "event": {
            "original": """{"event":{"transaction_id":1,"data_set":"my-logging-app"},"message":"WARN: Unable to get an interesting response"}""",
            "data_set": "my-logging-app",
            "transaction_id": 1
          }
        }
      }
    ]
  }
}

L'identifiant est spécifié de différentes manières selon l'outil, comme nous l'aborderons plus en détail dans les sections suivantes.

Beats

Pour les documents JSON, format courant pour de nombreuses sources de logs, si votre événement possède un identifiant utile et significatif pouvant servir d'identifiant unique pour un document et empêcher les entrées en double, l'utilisation du processeur decode_json_fields ou de json.document_ID paramètre d'entrée tel que recommandé dans la documentation. Cette approche est préférable à la génération d'une clé lorsqu'une clé naturelle est présente dans un champ JSON au sein de notre message. Les deux paramètres sont illustrés dans l'exemple ci-dessous :

filebeat.inputs:
- type: filestream
  id: my-logging-app
  paths:
    - /var/tmp/other.log
    - /var/log/*.log
  json.document_id: "event.transaction_id"  
# Alternative approach using decode_json_fields processor
processors:
  - decode_json_fields:
      document_id: "event.transaction_id"
      fields: ["message"]
      max_depth: 1
      target: ""

Pipelines d'ingestion

Dans ce cas, l'ID peut être défini à l'aide d'un processeur set combiné avec l'option copy_from pour transférer la valeur de votre champ unique vers le @metadata._id d'Elasticsearch attribut:

PUT _ingest/pipeline/test-pipeline
{
  "processors": [
    {
      "set": {
        "field": "_id",
        "copy_from": "transaction_id"
      }
    }
  ]
}

Elastic Agent

Elastic Agent adopte une approche similaire où vous pouvez utiliser le processeur copy_fields pour copier la valeur vers @metadata._id attribut dans l'intégration :

- copy_fields:
      fields:
        - from: transaction_id
          to: @metadata._id
      fail_on_error: true
      ignore_missing: true

Lorsque le paramètre fail_on_error est défini sur true, il entraîne un retour à l’état antérieur en annulant les modifications appliquées par le processeur en échec. Parallèlement, ignore_missing ne déclenchera une erreur pour un document comportant un champ inexistant que s’il est défini sur false.

ID auto-généré

Génération d’un ID unique à l’aide de techniques telles que l’empreinte numérique sur un sous-ensemble de champs d’événement. Le hachage d’un ensemble de champs génère une valeur unique qui, lorsqu’elle correspond, entraîne une mise à jour du document original lors de l’ingestion dans Elasticsearch. Comme l’explique cet article pratique sur la gestion des doublons avec Logstash décrit spécifiquement, le plug-in de filtre d’empreinte peut être configuré pour générer un ID avec l’algorithme de hachage spécifié dans le champ @metadata.fingerprint :

filter {
  fingerprint {
    source => ["event.start_date", "event.data_set", "message"]
    target => "[@metadata][fingerprint]"
    method => "SHA256"
  }
}
output {
  elasticsearch {
    hosts => "my-elastic-cluster.com"
    document_id => "%{[@metadata][fingerprint]}"
  }
}

Si rien n’est spécifié, l’algorithme de hachage par défaut SHA256 sera utilisé pour hacher la combinaison |event.start_date|start_date_value|event.data_set|data_set_value|message|message_value|. Si nous souhaitions utiliser l’une des autres options d’algorithme autorisées, il est possible de la spécifier à l’aide de l’option method. Cela aura pour résultat qu’Elasticsearch mettra à jour le document correspondant au _id généré :

GET my-fingerprinted-logstash-index/_search
{
  "took": 8,
  "timed_out": false,
  "_shards": {
    "total": 1,
    "successful": 1,
    "skipped": 0,
    "failed": 0
  },
  "hits": {
    "total": {
      "value": 1,
      "relation": "eq"
    },
    "max_score": 1,
    "hits": [
      {
        "_index": "my-fingerprinted-logstash-index",
        "_id": "b2faceea91b83a610bf64ac2b12e3d3b95527dc229118d8f819cdfaa4ba98af1",
        "_score": 1,
        "_source": {
          "@timestamp": "2023-10-23T16:46:00.772480Z",
          "message": "WARN: Unable to get an interesting response",
          "@version": "1",
          "event": {
            "original": """{"event":{"transaction_id":1,"data_set":"my-logging-app"},"message":"WARN: Unable to get an interesting response"}""",
            "data_set": "my-logging-app",
            "transaction_id": 1
          }
        }
      }
    ]
  }
}

Si votre événement ne possède pas de champ d'identification unique et significatif, cette option peut s'avérer utile si vous acceptez la surcharge de traitement liée à la génération de l'ID, ou le risque de collisions où différents événements aboutissent au même hachage généré. Des fonctionnalités similaires sont disponibles pour d'autres outils, comme nous l'aborderons dans les sections suivantes.

Beats

Le processeur add_id pour Beats et Elastic Agent permettra de générer un ID unique compatible avec Elasticsearch. Par défaut, cette valeur sera stockée dans @metadata._id champ qui est le champ d'ID pour les documents Elasticsearch.

filebeat.inputs:
- type: filestream
  ID: my-logging-app
  paths:
    - /var/tmp/other.log
    - /var/log/*.log
  json.document_ID: "event.transaction_id"  
processors:
  - add_ID: ~
      target_field: @metadata._id

Alternativement, le processeur fingerprint génère une valeur hachée d'une concaténation des paires de noms de champ et de valeurs spécifiées, séparées par l'opérateur |.

filebeat.inputs:
- type: filestream
  ID: my-logging-app
  paths:
    - /var/tmp/other.log
    - /var/log/*.log
processors:
  - fingerprint:
      fields: ["event.start_date", "event.data_set", "message"]
      target_field: "@metadata._id"
      method: "sha256"
      ignore_missing: false

Dans l'exemple ci-dessus, l'algorithme de hachage par défaut sha256 sera utilisé pour hacher la combinaison |event.start_date|start_date_value|event.data_set|data_set_value|message|message_value|. Si nous souhaitions utiliser l'une des autres options d'algorithme autorisées, elle peut être spécifiée à l'aide de l'option method. La gestion des erreurs est également un aspect important pour lequel l'option ignore_missing est utile. Par exemple, si le event.start_date le champ n'existe pas dans un document donné, une erreur sera générée lorsque ignore_missing est défini sur false. Il s'agit de l'implémentation par défaut si ignore_missing n'est pas explicitement défini, mais il est courant d'ignorer les erreurs en spécifiant la valeur true.

Elastic Agent

Tout comme Beats, Elastic Agent dispose d'un add_id processor qui peut être utilisé pour générer un identifiant unique, avec @metadata._id par défaut si l'attribut target_field n'est pas spécifié :

  - add_id:
      target_field: "@metadata._id"

Alternativement, le processeur fingerprint est également disponible dans Elastic Agent et peut être appliqué à tout segment d'intégration incluant une section de configuration avancée comprenant une option de processeurs. La logique du processeur se présente comme suit :

  - fingerprint:
      fields: ["event.start_date", "event.data_set", "message"]
      target_field: "@metadata._id"
      ignore_missing: false
      method: "sha256"

En prenant Kafka intégration à titre d'exemple, l'extrait de processeur ci-dessus peut être appliqué dans le segment de processeur de la section de configuration avancée pour Collect logs from Kafka brokers :

 

![Elastic Agent Kafka Integration Fingerprint Processor](./images/2.png)

 

Tout comme pour Beats, la valeur hachée est construite par concaténation du nom du champ et de la valeur du champ, séparés par |. Par exemple |field1|value1|field2|value2|. Cependant, tout comme dans Beats et contrairement à Logstash, la valeur method est en minuscules malgré la prise en charge des mêmes algorithmes d’encodage.

Pipelines d'ingestion

Nous allons présenter ici un exemple de requête pour créer un pipeline avec un fingerprint processor via l'API _ingest. Notez les similitudes de la configuration ci-dessous avec nos processeurs Beats :

PUT _ingest/pipeline/my-logging-app-pipeline
{
  "description": "Event and field dropping for my-logging-app",
  "processors": [
    {
      "fingerprint": {
        fields: ["event.start_date", "event.data_set", "message"]
        target_field: "@metadata._id"
        ignore_missing: false
        method: "SHA-256"
      }
    }
  ]
}

Agrégation d’événements

L’agrégation d’événements basée sur des champs communs est une autre option, si l’outil utilisé la prend en charge. L’agrégation d’événements implique un compromis, car l’outil doit conserver plusieurs événements en mémoire pour effectuer l’agrégation plutôt que de transférer immédiatement l’événement vers la sortie. Pour cette raison, le seul outil au sein de l’écosystème Elastic qui prend en charge l’agrégation d’événements est Logstash.

Pour implémenter l’approche basée sur l’agrégation dans Logstash, utilisez le plug-in aggregate. Dans notre cas, il est peu probable qu’un événement de fin spécifique soit envoyé pour distinguer les doublons ; il est donc nécessaire de spécifier un timeout, comme dans l’exemple ci-dessous, pour contrôler le processus de traitement par lots :

filter {
  grok {
    match => [ "message", %{NOTSPACE:event.start_date} "%{LOGLEVEL:loglevel} - %{NOTSPACE:user_ID} - %{GREEDYDATA:message}" ]
  }
  aggregate {
    task_ID => "%{event.start_date}%{loglevel}%{user_ID}"
    code => "map['error_count'] ||= 0; map['error_count'] += 1;"
    push_map_as_event_on_timeout => true
    timeout_task_ID_field => "user_id"
    timeout => 600
    timeout_tags => ['_aggregatetimeout']
    timeout_code => "event.set('has_multiple_occurrences', event.get('error_count') > 1)"
  }
}

L'exemple ci-dessus enverra un événement après 600 secondes, soit 10 minutes, en ajoutant les attributs error_count et has_multiple_occurrences à l'événement pour indiquer qu'il s'agit d'un événement agrégé. L'option push_map_as_event_on_timeout garantira que le résultat de l'agrégation est envoyé à chaque expiration de délai, ce qui vous permettra de réduire le volume d'alertes. Lors de la détermination du délai d'attente pour vos données, tenez compte de votre volume et optez pour le délai d'attente le plus court possible, car Logstash conservera les événements en mémoire jusqu'à ce que le délai expire et que l'événement agrégé soit envoyé.

Conclusions

Log volume spikes can quickly overwhelm logging platforms and SRE engineers looking to maintain reliable applications. We have discussed several approaches to handling duplicate events using Elastic Beats, Logstash (which are available in this GitHub repository), and Elastic Agent.

When generating IDs via a hashing algorithm using fingerprint processors, or performing aggregates, consider the attributes used carefully to balance preventing a flood and obfuscating legitimate streams pointing to a large-scale problem in your ecosystem. Both approaches have an overhead, either in terms of processing to generate the ID, or memory overhead to store the documents eligible to aggregate.

Selecting an option really depends on the events you consider duplicates and the performance trade-offs. As discussed, when you specify an ID Elasticsearch needs to check for the existence of a document matching that ID before adding the document to the index. This results in a slight delay in ingestion to perform the _id existence check.

Using hashing algorithms to generate the ID adds additional processing time as the ID needs to be generated for each event before it is compared and potentially ingested. Choosing to not specify an ID bypasses this check as Elastic will generate the ID for you, but will result in all events being stored which increases your storage footprint.

Dropping full events is a legitimate practice not covered in this piece. If you want to drop log entries to reduce your volume check out this piece on pruning fields from incoming log events.

If your favorite way to deduplicate events is not listed here, do let us know!

Ressources

  1. Elastic Beats
  2. Filebeat | Processeur Fingerprint
  3. Filebeat | Processeur de décodage de champs JSON
  4. Filebeat | Déduplication Filebeat
  5. Logstash
  6. Logstash | Plug-in Fingerprint
  7. Logstash | Plug-in Aggregate
  8. Elastic Agent
  9. Elastic Agent | Processeur Fingerprint
  10. Pipelines d’ingestion
  11. Elasticsearch | Processeur Fingerprint de pipeline d’ingestion