Log-Deduplizierung mit Elasticsearch

Doppelte Ereignisse von fehlerhaften Anwendungsdiensten erschweren die Logsuche. Erfahren Sie, wie Sie Duplikate mit Logstash, Beats und Elastic Agent verarbeiten.

Log-Deduplizierung mit Elasticsearch

SREs werden täglich von großen Mengen an Log aus unruhigen Anwendungen überflutet. In seinem wegweisenden Werk The Mythical Man Month, Frederick P. Brooks sagte: „Alle Programmierer sind Optimisten.“ Dieser Optimismus zeigt sich darin, dass Softwareentwickler keine Kontrollmechanismen implementieren, um ihre Anwendungen daran zu hindern, in außergewöhnlichen Fehlersituationen kontinuierlich Logdaten zu senden. In großen Unternehmen mit zentralisierten Logging-Plattformen wird diese Ereignisflut in Logging-Plattformen ingestiert und beansprucht beträchtliche Speichervolumen und Rechenleistung. Auf der menschlichen Seite führt dies dazu, dass sich SREs überfordert fühlen und unter Alerting-Müdigkeit leiden, da sie von der Nachrichtenwelle verschlungen werden, ein wenig wie hier:

 

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

 

Entwickler, die Software einschließlich Microservices und Schlüsselanwendungen erstellen, sind dafür verantwortlich, sicherzustellen, dass sie keine doppelten Log-Ereignisse senden und dass sie die korrekten Log-Ereignisse auf der richtigen Ebene senden. Dennoch bedeuten Situationen wie die Nutzung von solutions von Drittanbietern oder die Wartung veralteter Dienste, dass wir nicht immer garantieren können, dass verantwortungsvolle logging-Praktiken angewendet wurden. Selbst wenn wir unnötige Felder entfernen, wie in diesem Artikel über das Bereinigen von Feldern aus eingehenden Log-Ereignissen beschrieben, haben wir weiterhin ein Problem mit dem Speichern großer Mengen doppelter Ereignisse. Hier diskutieren wir die Herausforderungen bei der Identifizierung doppelter Logdaten aus problematischen Diensten und wie man Daten mithilfe von Elastic Beats, Logstash und Elastic Agent dedupliziert.

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.

Überblick über die Tools

In diesem Blog untersuchen wir die in vier Elastic-Tools verfügbaren Werkzeuge:

  1. Logstash ist ein kostenloses Open-Source-ETL-Pipeline-Tool, mit dem Sie Daten zwischen einer Vielzahl von Quellen ingestieren, transformieren und ausgeben können, einschließlich der Ingestion in und der Ausgabe aus Elasticsearch. Diese Beispiele werden verwendet, um das unterschiedliche Verhalten von Elasticsearch beim Generieren von Logdaten mit und ohne eine bestimmte ID zu zeigen.
  2. Beats sind eine Familie von leichtgewichtigen Shippern, die es uns ermöglichen, Ereignisse aus einer bestimmten Quelle nicht nur in Elasticsearch, sondern auch in andere Ausgänge zu ingestieren, einschließlich Kafka, Redis oder Logstash.
  3. Ingest-Pipelines ermöglichen die Anwendung von Transformationen und Anreicherungen auf Dokumente, die in Elasticsearch ingestiert wurden. Es ist, als würde man den filter-Teil von Logstash direkt in Elasticsearch ausführen, ohne dass ein weiterer Dienst laufen muss. Neue Pipelines können entweder über den Bildschirm Stack Management > Ingest Pipelines oder über die _ingest-API erstellt werden, wie in der Dokumentation beschrieben.
  4. Elastic Agent ist ein einzelner Agent, der auf Ihrem Host ausgeführt werden kann und Logdaten, Metriken sowie Sicherheitsdaten von mehreren Diensten und Infrastrukturen unter Verwendung der verschiedenen unterstützten Integrationen an Elasticsearch sendet. Unabhängig vom Grund für die Duplikate gibt es im Elastic-Ökosystem mehrere mögliche Vorgehensweisen.

Ingestion ohne angegebene IDs

Der Standardansatz besteht darin, alle Ereignisse zu ignorieren und zu ingestieren. Wenn für ein Dokument keine ID angegeben ist, generiert Elasticsearch automatisch eine neue ID für jedes empfangene Dokument. Wir nehmen ein einfaches Beispiel, das in diesem GitHub-Repository verfügbar ist und einen einfachen Express-HTTP-Server verwendet. Der Server stellt beim Ausführen einen einzelnen Endpoint bereit, der eine einzelne Log-Nachricht zurückgibt:

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

Mit Logstash können wir den Endpoint http://locahost:3000/ alle 60 Sekunden abfragen und das Ergebnis an Elasticsearch senden. Unsere logstash.conf sieht wie folgt aus:

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 überträgt jedes Ereignis, und ohne eine ID für das Ereignis generiert Elasticsearch eine neue _id Feld, das als eindeutige Kennung für jedes Dokument dient:

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
    ]
  }
}

Dieses Verhalten ist konsistent für Beats, Ingest-Pipelines und Elastic Agent, da sie alle empfangenen Ereignisse ohne zusätzliche Konfiguration senden.

BYO-ID

Die Angabe einer eindeutigen ID für jedes Ereignis unter Verwendung einer vorhandenen ID umgeht den im vorherigen Abschnitt besprochenen Schritt der Elasticsearch-ID-Generierung. Der Ingest eines Dokuments, bei dem dieses Attribut bereits existiert, führt dazu, dass Elasticsearch prüft, ob ein Dokument mit dieser ID im Index vorhanden ist, und das Dokument gegebenenfalls aktualisiert. Dies führt zu einem Mehraufwand, da der Index durchsucht werden muss, um zu prüfen, ob ein Dokument mit derselben _id bereits existiert. Erweiternd zu unserem obigen Logstash-Beispiel lässt sich der Wert der Dokument-ID in Logstash durch Angabe der Option document_id im Elasticsearch-Ausgang-Plugin festlegen, das für den Ingest von Ereignissen in Elasticsearch verwendet wird:

# 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]}"
    }
}

Dadurch wird der Wert des Feldes _id auf den Wert von event.transaction_id gesetzt. In unserem Ticket bedeutet dies, dass das neue Dokument bei der Ingestion das vorhandene Dokument ersetzt, da beide Dokumente eine _id von 1 haben:

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
          }
        }
      }
    ]
  }
}

Die ID wird je nach Tool auf unterschiedliche Weise angegeben, wie in den folgenden Abschnitten näher erläutert.

Beats

Wenn Ihr Ereignis für JSON-Dokumente, ein gängiges Format für viele Logquellen, über eine nützliche und aussagekräftige Kennung verfügt, die als eindeutige ID für ein Dokument verwendet werden kann und doppelte Einträge verhindert, können Sie entweder den Prozessor decode_json_fields oder die Einstellung json.document_ID verwenden Eingangseinstellung wie in der Dokumentation empfohlen. Dieser Ansatz wird der Generierung eines Schlüssels vorgezogen, wenn ein natürlicher Schlüssel innerhalb eines JSON-Felds in unserer Nachricht vorhanden ist. Beide Einstellungen sind im folgenden Beispiel dargestellt:

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: ""

Ingest-Pipelines

In diesem Fall kann die ID mithilfe eines set-Processors festgelegt werden in Kombination mit der Option copy_from, um den Wert aus Ihrem eindeutigen Feld an das Elasticsearch-Feld @metadata._id zu übertragen Attribut:

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

Elastic Agent

Elastic Agent verfolgt einen ähnlichen Ansatz, bei dem Sie den copy_fields-Prozessor verwenden können, um den Wert in @metadata._id zu kopieren. Attribut in der Integration:

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

Wenn die Einstellung fail_on_error auf „true“ gesetzt ist, wird der vorherige Zustand wiederhergestellt, indem die vom fehlerhaften Prozessor angewendeten Änderungen rückgängig gemacht werden. Unterdessen löst ignore_missing nur dann einen Fehler für ein Dokument mit einem nicht existierenden Feld aus, wenn es auf false gesetzt ist.

Automatisch generierte ID

Generierung einer eindeutigen ID unter Verwendung von Techniken wie Fingerprinting für eine Teilmenge von Ereignisfeldern. Durch das Hashing einer Reihe von Feldern wird ein eindeutiger Wert generiert, der bei einer Übereinstimmung zu einer Aktualisierung des ursprünglichen Dokuments beim Ingest in Elasticsearch führt. Wie dieser hilfreiche Artikel zum Umgang mit Duplikaten in Logstash wird spezifisch dargelegt, dass das Fingerprint Filter Plugin so konfiguriert werden kann, dass es eine ID mit dem angegebenen Hashing-Algorithmus für das Feld @metadata.fingerprint generiert:

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]}"
  }
}

Falls nicht anders angegeben, wird der Standard-Hashing-Algorithmus SHA256 verwendet, um die Kombination |event.start_date|start_date_value|event.data_set|data_set_value|message|message_value| zu hashen. Wenn wir eine der anderen zulässigen Algorithmus-Optionen verwenden wollten, kann dies über die Option method angegeben werden. Dies führt dazu, dass Elasticsearch das Dokument aktualisiert, das mit der generierten _id übereinstimmt:

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
          }
        }
      }
    ]
  }
}

Falls Ihr Ereignis kein einzelnes aussagekräftiges Identifizierungsfeld besitzt, kann dies eine nützliche Option sein, sofern Sie den Verarbeitungsaufwand für die Generierung der ID oder das Potenzial für Kollisionen in Kauf nehmen, bei denen unterschiedliche Ereignisse denselben generierten Hashwert ergeben. Ähnliche Funktionen sind für andere Tools verfügbar, wie in den nachfolgenden Abschnitten erläutert.

Beats

Der add_id-Prozessor für Beats und Elastic Agent ermöglicht die Generierung einer eindeutigen, mit Elasticsearch kompatiblen ID. Standardmäßig wird dieser Wert in @metadata._id gespeichert. Feld, das als ID-Feld für Elasticsearch-Dokumente dient.

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

Alternativ kann der fingerprint-Prozessor generiert einen Hash-Wert aus einer Verkettung der angegebenen Feldnamen- und Wertepaare, die durch den Operator | getrennt sind.

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

Im obigen Beispiel wird der Standard-Hashing-Algorithmus sha256 verwendet, um die Kombination |event.start_date|start_date_value|event.data_set|data_set_value|message|message_value| zu hashen. Wenn wir eine der anderen zulässigen Algorithmusoptionen verwenden möchten, kann dies über die Option method angegeben werden. Die Fehlerbehandlung ist ebenfalls ein wichtiger Aspekt, bei dem die Option ignore_missing hilfreich ist. Wenn zum Beispiel event.start_date Feld auf einem gegebenen Dokument nicht existiert, wird ein Fehler ausgelöst, wenn ignore_missing auf false gesetzt ist. Dies ist die Standardimplementierung, wenn ignore_missing nicht explizit festgelegt ist; es ist jedoch üblich, Fehler zu ignorieren, indem der Wert auf true gesetzt wird.

Elastic Agent

Genau wie Beats verfügt Elastic Agent über einen add_id-Prozessor, der zum Generieren einer eindeutigen ID verwendet werden kann. Standardmäßig wird @metadata._id verwendet, wenn das Attribut target_field nicht angegeben ist:

  - add_id:
      target_field: "@metadata._id"

Alternativ ist der fingerprint-Prozessor auch im Elastic Agent verfügbar und kann auf jedes Integrationssegment angewendet werden, das einen Abschnitt für erweiterte Konfigurationen einschließlich einer Prozessoroption enthält. Die Prozessorlogik sieht wie folgt aus:

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

Wenn man Kafka nimmt, Integration als Beispiel: Der obige Prozessor-Schnipsel kann im Prozessor-Segment des Abschnitts für die erweiterte Konfiguration unter Collect log from Kafka Broker angewendet werden:

 

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

 

Genau wie bei Beats wird der zu hashende Wert als Verkettung des Feldnamens und des Feldwerts konstruiert, die durch „|“ getrennt sind. Zum Beispiel |field1|value1|field2|value2|. Genau wie bei Beats und im Gegensatz zu Logstash ist der Wert für „method“ jedoch kleingeschrieben, obwohl dieselben Codierungsalgorithmen unterstützt werden.

Ingest-Pipelines

Hier zeigen wir die Beispielanfrage zum Erstellen einer Pipeline mit einem fingerprint-Prozessor über die _ingest-API. Beachten Sie die Ähnlichkeiten der folgenden Konfiguration mit unseren Beats-Prozessoren:

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"
      }
    }
  ]
}

Ereignisse aggregieren

Das Aggregieren von Ereignissen basierend auf gemeinsamen Feldern ist eine weitere Option, sofern das verwendete Tool dies unterstützt. Die Aggregation von Ereignissen ist mit einem Kompromiss verbunden, da das Tool mehrere Ereignisse im Arbeitsspeicher halten muss, um die Aggregation durchzuführen, anstatt das Ereignis sofort an den Ausgang weiterzuleiten. Aus diesem Grund ist Logstash das einzige Tool innerhalb des Elastic-Ökosystems, das Ereignisaggregation unterstützt.

Um den auf Aggregation basierenden Ansatz in Logstash zu implementieren, verwenden Sie das aggregate-Plugin. In unserem Ticket ist es unwahrscheinlich, dass ein spezifisches Endereignis gesendet wird, um Duplikate zu unterscheiden. Daher ist die Angabe eines timeout gemäß dem folgenden Beispiel erforderlich, um den Batch-Prozess zu steuern:

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)"
  }
}

Das obige Beispiel sendet nach 600 Sekunden bzw. 10 Minuten ein Ereignis und fügt dem Ereignis die Attribute error_count und has_multiple_occurrences hinzu, um ein aggregiertes Ereignis zu kennzeichnen. Die Option push_map_as_event_on_timeout stellt sicher, dass das Aggregationsergebnis bei jedem Timeout übertragen wird, wodurch Sie das Alert-Volumen reduzieren können. Berücksichtigen Sie bei der Festlegung des Timeouts für Ihre Daten das Volumen und wählen Sie das niedrigstmögliche Timeout, da Logstash die Ereignisse so lange im Arbeitsspeicher hält, bis das Timeout abläuft und das aggregierte Ereignis übertragen wird.

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!

Ressourcen

  1. Elastic Beats
  2. Filebeat | Fingerprint-Prozessor
  3. Filebeat | Prozessor zum Dekodieren von JSON-Feldern
  4. Filebeat | Filebeat-Deduplizierung
  5. Logstash
  6. Logstash | Fingerprint-Plugin
  7. Logstash | Aggregate-Plugin
  8. Elastic Agent
  9. Elastic Agent | Fingerprint-Prozessor
  10. Ingest-Pipelines
  11. Elasticsearch | Ingest-Pipeline-Fingerprint-Prozessor