# Replaying Dead Letter messages via RabbitMQ

**URL:** <https://discuss.elastic.co/t/replaying-dead-letter-messages-via-rabbitmq/137475>\
**Category:** Logstash\
**Created:** [June 26, 2018, 4:54pm UTC](https://discuss.elastic.co/t/replaying-dead-letter-messages-via-rabbitmq/137475 "2018-06-26T16:54:21Z")\
**Posts on this page:** 2\
**Page:** 1

<div class="post-metadata">

**Author:** ![davidbugeja](https://avatars.discourse-cdn.com/v4/letter/d/a5b964/32.png) [@davidbugeja](https://discuss.elastic.co/u/davidbugeja)\
**Post date:** [June 26, 2018, 4:54pm UTC](https://discuss.elastic.co/t/replaying-dead-letter-messages-via-rabbitmq/137475/1 "2018-06-26T16:54:21Z")

</div>

Logstash and Elasticion version: 5.6.9

Setting up multiple pipelines so that any messages that end up in Logstash's Dead Letter get sent back to RabbitMQ so that we can re-process them after sorting out any potential mapping issues (Further to [Logstash Dead Letter Messages back to RabbitMQ](https://discuss.elastic.co/t/logstash-dead-letter-messages-back-to-rabbitmq/136270))

Defined the following pipelines in Logstash:

```
########### Actual pipeline sending to Elastic ###########

input {
  rabbitmq {
    host => "rabbitmq"
    port => 5672
    user => "test"
    password => "test"
    vhost => "test"
    queue => "test.queue"
    codec => "json"
    metadata_enabled => true
    passive => true
    ack => true
    prefetch_count => 1000
    threads => 5
    tags => "audit"
    heartbeat => "5"
    }
}

filter {
  if "audit" in [tags] {
    mutate {
      rename => { "messageTemplate" => "message" }
      remove_field => ["@version"]
      remove_field => ["[fields][PerformanceType]" ]
   }
    if ![target_index] {
      mutate {
        add_field => { "target_index" => "%{[@metadata][rabbitmq_headers][x-target-index]}" }
      }
    }
  }
}

output {
  if "audit" in [tags] {
     elasticsearch {
        hosts => ["elastic01","elastic02","elastic03"]
        index => "%{target_index}"
        template => "/etc/logstash/templates/template1"
        template_name => "audit"
     }
  }
}

####### Dead Letter Pipeline to send to RabbitMQ #########

input {
  dead_letter_queue {
    path => "/opt/logstash/dead_letter_queue"
    commit_offsets => true
    sincedb_path => "/opt/logstash/dead_letter_queue/sincedb"
    tags => "dlx"
  }
}

filter {
  if "dlx" in [tags] {
        if ![rabbitmq.dlx.retry.count] {
                mutate { add_field => {"rabbitmq.dlx.retry.count" => "0" } }
                mutate { convert => { "rabbitmq.dlx.retry.count" => "integer" } }
        } else {
                ruby { code => 'event.set("rabbitmq.dlx.retry.count", event.get("rabbitmq.dlx.retry.count").to_i + 1)' }
        }
   if ![rabbitmq.dlx.reason] {
      mutate { add_field => { "rabbitmq.dlx.reason" => "%{[@metadata][dead_letter_queue][reason]}" } }
   } else {
      mutate { update => { "rabbitmq.dlx.reason" => "%{[@metadata][dead_letter_queue][reason]}" } }
   }
  }
}

output {
  if "dlx" in [tags] {
  rabbitmq {
    host => "rabbitmq"
    port => 5672
    user => "test"
    password => "test"
    vhost => "test"
    exchange => "test.dlx.exchange"
    exchange_type => "fanout"
    codec => "json"
    heartbeat => "5"
  }
  }
}

####### Healing Pipeline to sort out any mapping errors. No Output defined here as the removal of the dlx tag implies that the message goes through Output with audit in tags #########

input {
  rabbitmq {
    host => "rabbitmq"
    port => 5672
    user => "test"
    password => "test"
    vhost => "test"
    queue => "test.heal.queue"
    codec => "json"
    metadata_enabled => true
    passive => true
    ack => true
    prefetch_count => 1000
    threads => 5
    tags => "audit_heal"
    heartbeat => "5"
    }
}

filter {
  if "audit_heal" in [tags] {
    mutate {
      rename => { "messageTemplate" => "message" }
      remove_field => ["@version"]
      remove_field => ["[fields][PerformanceType]" ]
      remove_tag => ["dlx"]
   }
  }
}

```

In a nutshell, any messages that end up in Logstash's dead letter get sent to a special dlx queue in RabbitMQ with two further fields being a retry counter and the reason for dead letter. We then move the messages to a "heal" queue to process and re-attempt to send to Elastic.

My questions are:

1. If the message ends up in dead letter again after trying to fix it, will it be as a new document? Or will logstash consider it as an existing document in its dead letter and just leave it there?

2. Is there a way to forcefully create mapping errors in Logstash to cycle the same messages through logstash so I can test out the retry counter and see the dead letter reason field update on each try?

3. I noticed warnings in the logs `event previously submitted to dead letter queue. Skipping...` which I'm not sure why I'm getting as I have sincedb setup in the config?

Thanks in advance!

---

<div class="post-metadata">

**Author:** ![system](https://us1.discourse-cdn.com/elastic/original/3X/1/a/1ac57faf039f6b580b3f104ef42a2a89e41014de.png) [@system](https://discuss.elastic.co/u/system)\
**Post date:** [July 24, 2018, 4:54pm UTC](https://discuss.elastic.co/t/replaying-dead-letter-messages-via-rabbitmq/137475/2 "2018-07-24T16:54:26Z")

</div>

This topic was automatically closed 28 days after the last reply. New replies are no longer allowed.
