# Question about logstash and kafka

**URL:** <https://discuss.elastic.co/t/question-about-logstash-and-kafka/240553>\
**Category:** Logstash\
**Created:** [July 9, 2020, 3:26pm UTC](https://discuss.elastic.co/t/question-about-logstash-and-kafka/240553 "2020-07-09T15:26:52Z")\
**Posts on this page:** 14\
**Page:** 1

<div class="post-metadata">

**Author:** ![Maurizio\_Tarducci](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/maurizio_tarducci/32/71894_2.png) [@Maurizio\_Tarducci](https://discuss.elastic.co/u/Maurizio_Tarducci)\
**Post date:** [July 9, 2020, 3:26pm UTC](https://discuss.elastic.co/t/question-about-logstash-and-kafka/240553/1 "2020-07-09T15:26:52Z")

</div>

Hi i would like to configure logstash with kafka in a way that based on the content of the message that logstash collect send to kafka in a different topic : for example if the message contains the string XXX i want send the message to kafka in the topic XXX, if the message collected from logstach containg the string YYY i want send this message to kafka in the topic YYY.  
Is it possible configure logstash to do that ? if yes can you share with me an example of the config files ?  
Thanks in advance

---

<div class="post-metadata">

**Author:** ![Bouraoui\_Kacem](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/bouraoui_kacem/32/44744_2.png) [@Bouraoui\_Kacem](https://discuss.elastic.co/u/Bouraoui_Kacem)\
**Post date:** [July 9, 2020, 4:38pm UTC](https://discuss.elastic.co/t/question-about-logstash-and-kafka/240553/2 "2020-07-09T16:38:18Z")

</div>

Try this output

```auto
 kafka {
     bootstrap_servers => "kafka:9092"
     topic_id => "Topic name"
  }

```

---

<div class="post-metadata">

**Author:** ![Maurizio\_Tarducci](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/maurizio_tarducci/32/71894_2.png) [@Maurizio\_Tarducci](https://discuss.elastic.co/u/Maurizio_Tarducci)\
**Post date:** [July 9, 2020, 4:40pm UTC](https://discuss.elastic.co/t/question-about-logstash-and-kafka/240553/3 "2020-07-09T16:40:52Z")

</div>

Hi Bouraoui .. thanks this was my initial setting . I need something of more complicated like :

input {  
tcp {  
port =\> 5400  
codec =\> json  
}  
}  
output {  
if [message] =~ "XXX" {  
kafka {  
codec =\> json{}  
topic\_id =\> "XXX"  
}  
}  
else if [message] =~ "YYY" {  
kafka {  
codec =\> json{}  
topic\_id =\> "YYY"  
}  
}  
else {  
kafka {  
codec =\> json{}  
topic\_id =\> "ZZZ"  
}  
}  
}

---

<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:** [July 9, 2020, 4:51pm UTC](https://discuss.elastic.co/t/question-about-logstash-and-kafka/240553/4 "2020-07-09T16:51:34Z")

</div>

It is likely more efficient to do something like

```
filter
    if [message] =~ "XXX" {
        mutate { add_field => { "[@metadata][topic]" => "XXX" } }
    } else if [message] =~ "YYY" {
        mutate { add_field => { "[@metadata][topic]" => "YYY" } }
    } else {
        mutate { add_field => { "[@metadata][topic]" => "ZZZ" } }
    }
}
output {
    kafka {
        codec => json {}
        topic_id => "[@metadata][topic]"
    }
}
```

---

<div class="post-metadata">

**Author:** ![Maurizio\_Tarducci](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/maurizio_tarducci/32/71894_2.png) [@Maurizio\_Tarducci](https://discuss.elastic.co/u/Maurizio_Tarducci)\
**Post date:** [July 9, 2020, 4:53pm UTC](https://discuss.elastic.co/t/question-about-logstash-and-kafka/240553/5 "2020-07-09T16:53:10Z")

</div>

Ok thanks Badger  
a lot

---

<div class="post-metadata">

**Author:** ![Bouraoui\_Kacem](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/bouraoui_kacem/32/44744_2.png) [@Bouraoui\_Kacem](https://discuss.elastic.co/u/Bouraoui_Kacem)\
**Post date:** [July 9, 2020, 4:56pm UTC](https://discuss.elastic.co/t/question-about-logstash-and-kafka/240553/6 "2020-07-09T16:56:50Z")

</div>

Try to add a field to the message int the filter like this after extract the topicName :  
add\_field =\> {  
"[@metadata][topic]" =\> XXX  
}

Then you test in the output by :  
If([@metadata][topic]" == XXX) {

```
            }

```

I tried that with file type but not with kafka, i think it is the same

---

<div class="post-metadata">

**Author:** ![Maurizio\_Tarducci](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/maurizio_tarducci/32/71894_2.png) [@Maurizio\_Tarducci](https://discuss.elastic.co/u/Maurizio_Tarducci)\
**Post date:** [July 9, 2020, 5:18pm UTC](https://discuss.elastic.co/t/question-about-logstash-and-kafka/240553/7 "2020-07-09T17:18:39Z")

</div>

which is the variable that contains the complete event received from logstash ?  
i ask you this because it always match the 3rd condition .

---

<div class="post-metadata">

**Author:** ![Maurizio\_Tarducci](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/maurizio_tarducci/32/71894_2.png) [@Maurizio\_Tarducci](https://discuss.elastic.co/u/Maurizio_Tarducci)\
**Post date:** [July 9, 2020, 5:37pm UTC](https://discuss.elastic.co/t/question-about-logstash-and-kafka/240553/8 "2020-07-09T17:37:43Z")

</div>

This is an example in JSON format that arrive to Kafka ..

{"SessID":16843159,"APID":16843159,"Pgm":"JVMLDM80","Reads":1,"Type":"Directory","ACEELOGU":false,"JobNm":"TOMCAT","JSauth":false,"TokPRIV":false,"ACEEADSP":false,"TokRSPEC":false,"TokUDUS":false,"@version":"1","DirBlks":6,"Group":"OMVSGRP","Name":"####################","HostName":"XXXX","StepNm":"JAVAJVM","JobID":"S0012440","AGPID":16843159,"OUid":900005,"Open":"2020-07-09T17:15:54.298","ACEEAUDT":false,"@timestamp":"2020-07-09T17:15:54.442Z","DevNo":"00000a","Close":"2020-07-09T17:15:54.298","Inode":185,"Cat":"FS","ACEEROA":false,"Time":"2020-07-09T17:15:54.299","ASessID":16843159,"RType":92,"OGroup":1,"SubT":"File close","Token":7407136,"TokTRST":false,"ACEE":true,"PID":16843159,"FName":"/SND1/u/tomcat/conf/Catalina/snd1.bmc.com","SAF":1,"SAFD":"XXXX","PGroup":16843159,"ACEEPRIV":false,"host":"192.168.0.1","ACEESPEC":false,"TokSUS":false,"SessType":"Started Procedure","Start":"2020-07-05T23:07:23.230","ACEEFLG1":"Defined","Severity":"Info","port":12679,"TokFlg1":"Pre 1.9","WorkTypeD":"Started task","PrivStatD":"Normal user","BlksRead":13,"SID":"SND1","ACEEOPER":false,"TokFlg3":"Default SECLABEL, Default Group","UserID":"TOMCAT"}

---

<div class="post-metadata">

**Author:** ![Maurizio\_Tarducci](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/maurizio_tarducci/32/71894_2.png) [@Maurizio\_Tarducci](https://discuss.elastic.co/u/Maurizio_Tarducci)\
**Post date:** [July 9, 2020, 5:38pm UTC](https://discuss.elastic.co/t/question-about-logstash-and-kafka/240553/9 "2020-07-09T17:38:42Z")

</div>

the config suggested is not able to find a match in the above msg received as i expect .  
Thanks

---

<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:** [July 9, 2020, 5:49pm UTC](https://discuss.elastic.co/t/question-about-logstash-and-kafka/240553/10 "2020-07-09T17:49:02Z")

</div>

Are you using a json codec on the input? If you are then you will not have a field called [message]

---

<div class="post-metadata">

**Author:** ![Maurizio\_Tarducci](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/maurizio_tarducci/32/71894_2.png) [@Maurizio\_Tarducci](https://discuss.elastic.co/u/Maurizio_Tarducci)\
**Post date:** [July 9, 2020, 5:53pm UTC](https://discuss.elastic.co/t/question-about-logstash-and-kafka/240553/11 "2020-07-09T17:53:16Z")

</div>

yes i have json codec in input .  
input {  
tcp {  
port =\> 5400  
codec =\> json  
}  
}

---

<div class="post-metadata">

**Author:** ![Maurizio\_Tarducci](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/maurizio_tarducci/32/71894_2.png) [@Maurizio\_Tarducci](https://discuss.elastic.co/u/Maurizio_Tarducci)\
**Post date:** [July 9, 2020, 5:57pm UTC](https://discuss.elastic.co/t/question-about-logstash-and-kafka/240553/12 "2020-07-09T17:57:12Z")

</div>

which variable can i use ? instead of message ?

---

<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:** [July 9, 2020, 6:26pm UTC](https://discuss.elastic.co/t/question-about-logstash-and-kafka/240553/13 "2020-07-09T18:26:23Z")

</div>

If you want to check whether a particular field contains XXX/YYY/etc. then test that field. If you want to check whether _any_ field contains it, then change the codec to plain, test [message] to set [@metadata][topic] and then _after_ that use

```
 json { source => "message" remove_field => ["message"] }
```

---

<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:** [August 6, 2020, 6:34pm UTC](https://discuss.elastic.co/t/question-about-logstash-and-kafka/240553/14 "2020-08-06T18:34:35Z")

</div>

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