# Kafka input plugin - decorate\_events does not take effect

**URL:** <https://discuss.elastic.co/t/kafka-input-plugin-decorate-events-does-not-take-effect/136712>\
**Category:** Logstash\
**Created:** [June 20, 2018, 2:20pm UTC](https://discuss.elastic.co/t/kafka-input-plugin-decorate-events-does-not-take-effect/136712 "2018-06-20T14:20:44Z")\
**Posts on this page:** 13\
**Page:** 1

<div class="post-metadata">

**Author:** ![ArunANayagam](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/arunanayagam/32/45171_2.png) [@ArunANayagam](https://discuss.elastic.co/u/ArunANayagam)\
**Post date:** [June 20, 2018, 2:20pm UTC](https://discuss.elastic.co/t/kafka-input-plugin-decorate-events-does-not-take-effect/136712/1 "2018-06-20T14:20:44Z")

</div>

Hi,

I have a Kafka logstash input plugin that's reading messages from a Kafka 1.0.0 fine.

But the "decorate\_events" property does not seem to take effect, i.e, I receive none of the kafka topic/partition information, here's my input code.  
I am using Logstash/ES 6.2.4

```
input {
    kafka {
        bootstrap_servers => "kafka.dev:9092"
        topics => ["retail.clog.1", "digital.clog.1"]
        client_id => "log"
        group_id => "log"
        auto_offset_reset => "earliest"
        consumer_threads => 4
        decorate_events => true
        heartbeat_interval_ms => "3000"
        max_poll_interval_ms => "1000"
        reconnect_backoff_ms => "3000"
        #queue_size => 500
    }
}

```

Appreciate your help on this.

Thanks,  
Arun

---

<div class="post-metadata">

**Author:** ![Badger](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/badger/32/25190_2.png) [@Badger](https://discuss.elastic.co/u/Badger)\
**Post date:** [June 20, 2018, 2:25pm UTC](https://discuss.elastic.co/t/kafka-input-plugin-decorate-events-does-not-take-effect/136712/2 "2018-06-20T14:25:49Z")

</div>

Are you looking at [@metadata][kafka] in logstash? It will not get written to elasticsearch.

---

<div class="post-metadata">

**Author:** ![ArunANayagam](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/arunanayagam/32/45171_2.png) [@ArunANayagam](https://discuss.elastic.co/u/ArunANayagam)\
**Post date:** [June 20, 2018, 2:29pm UTC](https://discuss.elastic.co/t/kafka-input-plugin-decorate-events-does-not-take-effect/136712/3 "2018-06-20T14:29:35Z")

</div>

Hi,

Thanks for a quick response. Didn't quite understand, do I have to explicitly assign [@metadata][kafka] to a variable in logstash, if so, can you please say how.

Rest of the code is just this,

```
filter {
    json {
        source => "message"
    }
}
output {
	elasticsearch {
		action => "index"
		index => "ms-logs-%{+YYYY.MM.dd}"
		hosts => ["elastic.dev"]
	}
}

```

Thanks,  
Arun

---

<div class="post-metadata">

**Author:** ![Badger](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/badger/32/25190_2.png) [@Badger](https://discuss.elastic.co/u/Badger)\
**Post date:** [June 20, 2018, 2:36pm UTC](https://discuss.elastic.co/t/kafka-input-plugin-decorate-events-does-not-take-effect/136712/4 "2018-06-20T14:36:27Z")

</div>

The kafka input puts the topic, partition etc. into fields under [@metadata][kafka]. The [@metadata] field on the event exists in logstash, but is not written to elasticsearch, so if you want to have that data in elasticsearch you would need to use mutate+copy to copy [@metadata][kafka] to another field.

---

<div class="post-metadata">

**Author:** ![ArunANayagam](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/arunanayagam/32/45171_2.png) [@ArunANayagam](https://discuss.elastic.co/u/ArunANayagam)\
**Post date:** [June 20, 2018, 2:42pm UTC](https://discuss.elastic.co/t/kafka-input-plugin-decorate-events-does-not-take-effect/136712/5 "2018-06-20T14:42:15Z")

</div>

> [@Badger](#):
>
> [@metadata][kafka]

Thank you so much, that helped.

```
filter {
        json {
                source => "message"
        }
        mutate {
            add_field => {
                "kafka" => "%{[@metadata][kafka]}"
            }
        }
}

```

---

<div class="post-metadata">

**Author:** ![ArunANayagam](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/arunanayagam/32/45171_2.png) [@ArunANayagam](https://discuss.elastic.co/u/ArunANayagam)\
**Post date:** [June 20, 2018, 2:52pm UTC](https://discuss.elastic.co/t/kafka-input-plugin-decorate-events-does-not-take-effect/136712/6 "2018-06-20T14:52:26Z")

</div>

Hi,

Is there a way to index the kafka detail as well?  
I think because its an inner json it doesn't index in Elasticsearch.  
It would be nice to filter out messages of a particular partition in Kibana, for example.

Thanks,  
Arun

---

<div class="post-metadata">

**Author:** ![Badger](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/badger/32/25190_2.png) [@Badger](https://discuss.elastic.co/u/Badger)\
**Post date:** [June 20, 2018, 2:58pm UTC](https://discuss.elastic.co/t/kafka-input-plugin-decorate-events-does-not-take-effect/136712/7 "2018-06-20T14:58:12Z")

</div>

> [@ArunANayagam](#):
>
> mutate { add\_field =\> { "kafka" =\> "%{[@metadata][kafka]}" } }

mutate { copy =\> { "[@metadata][kafka]" =\> "kafka" } }

---

<div class="post-metadata">

**Author:** ![ArunANayagam](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/arunanayagam/32/45171_2.png) [@ArunANayagam](https://discuss.elastic.co/u/ArunANayagam)\
**Post date:** [June 20, 2018, 3:07pm UTC](https://discuss.elastic.co/t/kafka-input-plugin-decorate-events-does-not-take-effect/136712/8 "2018-06-20T15:07:19Z")

</div>

Hi,

I get this error now,

[2018-06-20T15:05:20,979][WARN][logstash.outputs.elasticsearch] Could not index event to Elasticsearch.  
{:status=\>400, :action=\>["index", {:\_id=\>nil, :\_index=\>"ms-logs-2018.06.20", :\_type=\>"doc", :\_routing=\>nil}, #LogStash::Event:0x354fa377],  
:response=\>{"index"=\>{"\_index"=\>"ms-logs-2018.06.20", "\_type"=\>"doc", "\_id"=\>"ua66HWQBIApBjZNHh1si", "status"=\>400, "error"=\>{"type"=\>"mapper\_parsing\_exception",  
"reason"=\>"failed to parse [kafka]", "caused\_by"=\>{"type"=\>"illegal\_state\_exception",  
"reason"=\>"Can't get text on a START\_OBJECT at 1:248"}}}}}

Thanks,  
Arun

---

<div class="post-metadata">

**Author:** ![Badger](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/badger/32/25190_2.png) [@Badger](https://discuss.elastic.co/u/Badger)\
**Post date:** [June 20, 2018, 3:12pm UTC](https://discuss.elastic.co/t/kafka-input-plugin-decorate-events-does-not-take-effect/136712/9 "2018-06-20T15:12:23Z")

</div>

Look at the elasticseach log. That will have a clearer error message.

---

<div class="post-metadata">

**Author:** ![ArunANayagam](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/arunanayagam/32/45171_2.png) [@ArunANayagam](https://discuss.elastic.co/u/ArunANayagam)\
**Post date:** [June 20, 2018, 3:18pm UTC](https://discuss.elastic.co/t/kafka-input-plugin-decorate-events-does-not-take-effect/136712/10 "2018-06-20T15:18:14Z")

</div>

Hi,

Here's the error from elastic logs

[2018-06-20T15:14:41,638][DEBUG][o.e.a.b.TransportShardBulkAction] [ms-logs-2018.06.20][2] failed to execute bulk item (index) BulkShardRequest [[ms-logs-2018.06.20][2]] containing [index {[ms-logs-2018.06.20][doc][CV\_DHWQBFf7QtJ8OFSFl], source[{"app":"log","@timestamp":"2018-06-20T15:14:41.528Z","latency":"340","status":"true","ts":"2018-06-20 15:14:41,523","lcruid":"300ac294-1997-432f-a4e2-9da04a5c9","iid":"732a0577b","@version":"1","lccid":"e3ab50c1-c984-4e5d-b954-b657bef57cf","kafka":{"offset":11,"timestamp":1529507681524,"partition":1,"topic":"digital.log.1","key":"e0bc3806-9421-41c2-a70a-744d11d5d","consumer\_group":"clog"},"type":"RESPONSE","msg":""Response from Service Layer"","crid":"arun1234","source":"DF"}]}]  
org.elasticsearch.index.mapper.MapperParsingException: failed to parse [kafka]  
at org.elasticsearch.index.mapper.FieldMapper.parse(FieldMapper.java:302) ~[elasticsearch-6.2.4.jar:6.2.4]

---

<div class="post-metadata">

**Author:** ![Badger](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/badger/32/25190_2.png) [@Badger](https://discuss.elastic.co/u/Badger)\
**Post date:** [June 20, 2018, 3:20pm UTC](https://discuss.elastic.co/t/kafka-input-plugin-decorate-events-does-not-take-effect/136712/11 "2018-06-20T15:20:53Z")

</div>

I suspect the problem is that you previously indexed documents where "kafka" was a string, and now it is an object. Are you in a position to "DELETE ms-logs-2018.06.20"? Or can you wait until tomorrow and see if it starts working at midnight UTC when it rolls to a new index?

---

<div class="post-metadata">

**Author:** ![ArunANayagam](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/arunanayagam/32/45171_2.png) [@ArunANayagam](https://discuss.elastic.co/u/ArunANayagam)\
**Post date:** [June 20, 2018, 3:24pm UTC](https://discuss.elastic.co/t/kafka-input-plugin-decorate-events-does-not-take-effect/136712/12 "2018-06-20T15:24:03Z")

</div>

Hi,

That was precisely the issue, no problem deleting the index, I am just setting this up.  
Thank you so much, you have been very helpful.

Thanks,  
Arun

---

<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 18, 2018, 3:24pm UTC](https://discuss.elastic.co/t/kafka-input-plugin-decorate-events-does-not-take-effect/136712/13 "2018-07-18T15:24:12Z")

</div>

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