Deduplicación de log con Elasticsearch
Los eventos duplicados de servicios de aplicaciones no saludables hacen que la búsqueda de log sea complicada. Consulta cómo manejar duplicados usando Logstash, Beats y Elastic Agent.

Deduplicación de log con Elasticsearch
Los SRE se ven inundados por grandes volúmenes de logs de aplicaciones ruidosas todos los días. En su obra fundamental The Mythical Man Month, Frederick P. Brooks dijo que "todos los programadores son optimistas". Este optimismo se manifiesta en que los ingenieros de software no implementan controles para evitar que sus aplicaciones envíen logs continuos en situaciones de fallo excepcionales. En organizaciones grandes con plataformas de logging centralizadas, esta avalancha de eventos se ingesta en las plataformas de logging y ocupa volúmenes de almacenamiento y capacidad de procesamiento considerables. En cuanto a las personas, esto hace que los SRE se sientan abrumados y sufran fatiga por alertas al verse envueltos por la ola de mensajes, algo así:
](https://images.contentstack.io/v3/assets/bltefdd0b53724fa2ce/bltcb726be86c9c8601/6543831208cc0104077cd7d4/gif.gif)
Los desarrolladores que crean software, incluidos microservicios y aplicaciones clave, son responsables de garantizar que no envíen eventos de log duplicados y que envíen los eventos de log correctos al nivel adecuado. Sin embargo, situaciones como el uso de soluciones de terceros o el mantenimiento de servicios antiguos significan que no siempre podemos garantizar que se hayan aplicado prácticas de logging responsables. Incluso si eliminamos los campos innecesarios como se explica en este artículo sobre la poda de campos de eventos de log entrantes, todavía tenemos un problema con almacenar grandes cantidades de eventos duplicados. Aquí analizamos los desafíos de identificar logs duplicados de servicios problemáticos y cómo deduplicar datos mediante Elastic Beats, Logstash y 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.
Visión general de las herramientas
En este blog examinaremos las herramientas disponibles en cuatro herramientas de Elastic:
- Logstash es una herramienta de pipeline ETL gratuita y open source que te permite ingestar, transformar y generar datos entre una miríada de fuentes, incluida la ingesta en Elasticsearch y la salida desde este. Estos ejemplos se usarán para mostrar el comportamiento diferente de Elasticsearch al generar logs con y sin un ID particular.
- Beats son una familia de agentes ligeros que nos permiten la ingesta de eventos desde una fuente determinada no solo en Elasticsearch, sino también en otras salidas, incluidas Kafka, Redis o Logstash.
- Pipelines de ingesta permiten aplicar transformaciones y enriquecimiento a los documentos ingestados en Elasticsearch. Es como ejecutar la parte de
filterde Logstash directamente en Elasticsearch sin necesidad de tener otro servicio en ejecución. Se pueden crear nuevos pipelines dentro de la pantalla Stack Management > Ingesta Pipelines o mediante la API_ingesta, como se explica en la documentación. - Elastic Agent es un agente único que puede ejecutarse en tu host y enviar logs, métricas y datos de seguridad desde múltiples servicios e infraestructura a Elasticsearch mediante las diversas integraciones soportadas. Independientemente del motivo de los duplicados, existen varios cursos de acción posibles en el ecosistema de Elastic.
Ingesta sin ID especificados
El enfoque predeterminado es ignorar e ingestar todos los eventos. Cuando no se especifica un ID en un documento, Elasticsearch generará automáticamente un nuevo ID para cada documento que reciba. Tomemos un ejemplo simple, disponible en este repositorio de GitHub, usando un servidor HTTP Express simple. El servidor, cuando se ejecuta, expone un único endpoint que devuelve un solo mensaje de log:
{"event":{"transaction_id":1,"data_set":"my-logging-app"},"message":"WARN: Unable to get an interesting response"}Usando Logstash podemos consultar el endpoint http://locahost:3000/ cada 60 segundos y enviar el resultado a Elasticsearch. Nuestro logstash.conf se ve como lo siguiente:
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 enviará cada evento y, sin ningún ID en el evento, Elasticsearch generará un nuevo _id campo que sirva como identificador único para cada documento:
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
]
}
}Este comportamiento es consistente para Beats, los pipelines de ingesta y Elastic Agent, ya que enviarán todos los eventos recibidos sin configuración adicional.
BYO-ID
Especificar un ID único para cada evento mediante un ID existente evita el paso de generación de ID de Elasticsearch que se trató en la sección anterior. La ingesta de un documento donde este atributo ya existe resultará en que Elasticsearch verifique si un documento con este ID existe en el índice y, de ser así, actualice el documento. Esto resulta en una sobrecarga, ya que es necesario realizar una búsqueda en el índice para comprobar si existe un documento con el mismo _id.
Ampliando nuestro ejemplo de Logstash anterior, es posible especificar el valor del ID de documento en Logstash mediante la especificación de la opción document_id en el plugin de salida de Elasticsearch, la cual se usaría para la ingesta de eventos en 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]}"
}
}Esto establecerá el valor del campo _id en el valor de event.transaction_id. En nuestro caso, esto significa que el nuevo documento reemplazará al documento existente durante la ingesta, ya que ambos documentos tienen 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
}
}
}
]
}
}El ID se especifica de varias maneras según la herramienta, como se analiza más adelante en las secciones siguientes.
Beats
Para documentos JSON, que es un formato común para muchas fuentes de logs, si tu evento tiene un ID útil y significativo que pueda usarse como ID único para un documento y prevenir entradas duplicadas, el uso del procesador decode_json_fields o json.document_ID configuración de entrada como se recomienda en la documentación. Este enfoque es preferible a generar una clave cuando hay una clave natural presente dentro de un campo JSON en nuestro mensaje.
Ambas configuraciones se muestran en el siguiente ejemplo:
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 de ingesta
En este caso, el ID se puede establecer mediante un set processor combinado con la opción copy_from para transferir el valor desde tu campo único al @metadata._id de Elasticsearch atributo:
PUT _ingest/pipeline/test-pipeline
{
"processors": [
{
"set": {
"field": "_id",
"copy_from": "transaction_id"
}
}
]
}Elastic Agent
Elastic Agent tiene un enfoque similar en el que puedes usar el procesador copy_fields para copiar el valor a @metadata._id atributo en la integración:
- copy_fields:
fields:
- from: transaction_id
to: @metadata._id
fail_on_error: true
ignore_missing: trueCuando la configuración fail_on_error es verdadera, resultará en un retorno al estado anterior al revertir los cambios aplicados por el procesador que falló. Mientras tanto, ignore_missing solo provocará un error para un documento con un campo inexistente cuando se configure como false.
ID generado automáticamente
Generación de un ID único mediante técnicas como fingerprinting en un subconjunto de campos de eventos. Al aplicar hash a un conjunto de campos, se genera un valor único que, al coincidir, resultará en una actualización del documento original durante la ingesta en Elasticsearch.
Como se explica en este útil artículo sobre el manejo de duplicados con Logstash describe específicamente que el plugin de filtro fingerprint puede configurarse para generar un ID con el algoritmo de hash especificado en el campo @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 no se especifica, se usará el algoritmo de hash predeterminado SHA256 para hashear la combinación |event.start_date|start_date_value|event.data_set|data_set_value|message|message_value|. Si quisiéramos usar una de las otras opciones de algoritmo permitidas, se puede especificar usando la opción method. Esto dará como resultado que Elasticsearch actualice el documento que coincide con el _id generado:
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 tu evento no tiene un único campo de identificación significativo, esta puede ser una opción útil si estás dispuesto a asumir la sobrecarga de procesamiento de generar la ID, o la posibilidad de colisiones donde diferentes eventos se resuelvan en el mismo hash generado. Hay capacidades similares disponibles para otras herramientas, como se analiza en las secciones posteriores.
Beats
El procesador add_id para Beats y Elastic Agent permitirá generar una ID única compatible con Elasticsearch. De forma predeterminada, este valor se almacenará en @metadata._id campo que es el campo de ID para documentos de 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._idAlternativamente, el procesador fingerprint genera un valor hash de una concatenación de los pares de nombre de campo y valor especificados separados por el operador |.
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: falseEn el ejemplo anterior, se utilizará el algoritmo de hashing predeterminado sha256 para hacer hash de la combinación |event.start_date|start_date_value|event.data_set|data_set_value|message|message_value|. Si quisiéramos usar una de las otras opciones de algoritmo permitidas, se puede especificar usando la opción method.
El manejo de errores también es una consideración importante con la que ayuda la opción ignore_missing. Por ejemplo, si event.start_date el campo no existe en un documento determinado, se generará un error cuando ignore_missing esté configurado como false. Esta es la implementación predeterminada si ignore_missing no se establece explícitamente, pero es común ignorar los errores especificando el valor como true.
Elastic Agent
Al igual que Beats, Elastic Agent tiene un procesador add_id que se puede usar para generar un ID único, que usa @metadata._id de forma predeterminada si no se especifica el atributo target_field:
- add_id:
target_field: "@metadata._id"Como alternativa, el fingerprint processor también está disponible en Elastic Agent y se puede aplicar a cualquier segmento de integración que incluya una sección de configuración avanzada con una opción processors. La lógica del procesador se ve así:
- fingerprint:
fields: ["event.start_date", "event.data_set", "message"]
target_field: "@metadata._id"
ignore_missing: false
method: "sha256"Tomando Kafka integración como ejemplo, el fragmento de procesador anterior se puede aplicar en el segmento de procesador de la sección de configuración avanzada para Collect logs from Kafka brokers (Recopilar logs de brokers de Kafka):
](https://images.contentstack.io/v3/assets/bltefdd0b53724fa2ce/blt6a63c6db6d13008d/654384290970dd001bd1637d/image_2_kafka.png)
Al igual que en Beats, el valor que se hashea se construye como una concatenación del nombre del campo y el valor del campo separados por |. Por ejemplo |field1|value1|field2|value2|. Sin embargo, al igual que en Beats y a diferencia de Logstash, el valor method está en minúsculas a pesar de admitir los mismos algoritmos de codificado.
Pipelines de ingesta
Aquí mostraremos la solicitud de ejemplo para crear un pipeline con un procesador fingerprint con la API _ingest. Observa las similitudes de la siguiente configuración con nuestros procesadores de 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"
}
}
]
}Agregar eventos
Agregar eventos en conjunto basados en campos comunes es otra opción, si la herramienta que utilizas brinda soporte para ello. La agregación de eventos conlleva una contrapartida, ya que la herramienta necesita mantener varios eventos en la memoria para realizar la agregación en lugar de reenviar el evento inmediatamente a la salida. Por esta razón, la única herramienta dentro del ecosistema de Elastic que brinda soporte para la agregación de eventos es Logstash.
Para implementar el enfoque basado en agregación en Logstash, usa el aggregate plugin. En nuestro caso, es poco probable que se envíe un evento final específico para distinguir entre duplicados, lo que significa que es necesario especificar un timeout según el ejemplo a continuación para controlar el proceso de batch:
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)"
}
}El ejemplo anterior enviará un evento después de 600 segundos, o 10 minutos, añadiendo los atributos error_count y has_multiple_occurrences al evento para indicar un evento agregado. La opción push_map_as_event_on_timeout garantizará que el resultado de la agregación se envíe en cada tiempo de espera, lo que te permitirá reducir el volumen de alertas. Al determinar el tiempo de espera para tus datos, considera tu volumen y opta por el tiempo de espera más bajo posible, ya que Logstash mantendrá los eventos en la memoria hasta que el tiempo de espera expire y se envíe el evento agregado.
Recursos de ## 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!
- Elastic Beats
- Filebeat | Procesador de fingerprint
- Filebeat | Procesador de decodificar campos JSON
- Filebeat | Deduplicación de Filebeat
- Logstash
- Logstash | Plugin de fingerprint
- Logstash | Plugin de aggregate
- Elastic Agent
- Elastic Agent | Procesador de fingerprint
- Pipelines de ingesta
- Elasticsearch | Procesador de fingerprint de pipeline de ingesta