# Logstash Kafka input : conditonnal consuming?

**URL:** <https://discuss.elastic.co/t/logstash-kafka-input-conditonnal-consuming/268664>\
**Category:** Logstash\
**Created:** [March 29, 2021, 2:08pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-conditonnal-consuming/268664 "2021-03-29T14:08:25Z")\
**Posts on this page:** 14\
**Page:** 1

<div class="post-metadata">

**Author:** ![Travis](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/travis/32/54079_2.png) [@Travis](https://discuss.elastic.co/u/Travis)\
**Post date:** [March 29, 2021, 2:08pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-conditonnal-consuming/268664/1 "2021-03-29T14:08:25Z")

</div>

Hello !

I have multiple firewall log sources. All these sources push the logs in a dedicated topic named "firewall".

In logstash, I would like to have a different pipeline for each of these sources to apply different processing and use different index. For example, one pipeline for Fortigate and one pipeline for Juniper.

Let's check with the following Fortigate pipeline (I didn't changed kafka group id which is by default "logstash") :

```
input {
  kafka {
   topics => ["firewall"]
   codec => json
   tags => ["Fortigate"]
}
}

filter{
}

output {
if "Fortigate" in [tags] {
  elasticsearch {
  hosts => ["elastic1:9200"]
 index => "firewall"
}
}
}

```

It works. BUT if I specify a different filter in my output as this :

```
output {
if "TEST" in [tags] {
  elasticsearch {
  hosts => ["elastic1:9200"]
}
}
}

```

I can see that topic is still consumed (lag is not increasing). For me it should not be consumed because the tag is not the good one in the output.

This way, If i create another pipeline "Juniper", I think that data will be already consumed by the Fortigate pipeline.

So, does it means that, **in any case, all data is consumed whatever is the filter used in my output ?**

How can I deal with my needs ?

Thanks for your help !

---

<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:** [March 29, 2021, 4:40pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-conditonnal-consuming/268664/2 "2021-03-29T16:40:32Z")

</div>

Not sure I understood the question, but if an output is conditional

```
output {
    if "Fortigate" in [tags] {
        elasticsearch {

```

if the condition does not evaluate to true then the event is not sent to an output, it is discarded.

---

<div class="post-metadata">

**Author:** ![Travis](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/travis/32/54079_2.png) [@Travis](https://discuss.elastic.co/u/Travis)\
**Post date:** [March 30, 2021, 2:50pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-conditonnal-consuming/268664/3 "2021-03-30T14:50:44Z")

</div>

Yes I agree with this.

My question is : Even if I don't use an output, logs will be consumed ?

---

<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:** [March 30, 2021, 4:20pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-conditonnal-consuming/268664/4 "2021-03-30T16:20:24Z")

</div>

If you do not define an output section then I do not think the pipeline will be executed, however, the output section does not have to send anything to an output. It is OK if the conditional is never true.

---

<div class="post-metadata">

**Author:** ![Travis](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/travis/32/54079_2.png) [@Travis](https://discuss.elastic.co/u/Travis)\
**Post date:** [March 30, 2021, 4:34pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-conditonnal-consuming/268664/5 "2021-03-30T16:34:42Z")

</div>

I made a test (no output,no filter, just my kafka input) and when looking at kafka metrics on a kafka machine (kafka-consumer-group.sh ...), I can see logs are consumed as lag is not increasing. If it was not consumed, this would not be the case (in my understanding of how kafka works)

So even if I don't have any output, pipeline is executed.

---

<div class="post-metadata">

**Author:** ![rcowart](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/rcowart/32/88091_2.png) [@rcowart](https://discuss.elastic.co/u/rcowart)\
**Post date:** [March 30, 2021, 8:03pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-conditonnal-consuming/268664/6 "2021-03-30T20:03:13Z")

</div>

I think that you misunderstand how Kafka works. If you want two different pipelines to both be able to consume all of the events in a topic, each pipeline must be configured for a separate consumer group using the `group_id` option.

---

<div class="post-metadata">

**Author:** ![Travis](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/travis/32/54079_2.png) [@Travis](https://discuss.elastic.co/u/Travis)\
**Post date:** [March 31, 2021, 6:47am UTC](https://discuss.elastic.co/t/logstash-kafka-input-conditonnal-consuming/268664/7 "2021-03-31T06:47:20Z")

</div>

Thank you Rob. It makes sense after reading more in depth Kafka documentation.

So, ok I use a different group id for each pipeline (cisco-pipeline.conf, juniper-pipeline.conf..).

But I'm asking what becomes the messages **that does not match my filter**?

Are they discarded once consumed/acknowledged ? They are never wrote ? even temporary in memory ?

---

<div class="post-metadata">

**Author:** ![rcowart](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/rcowart/32/88091_2.png) [@rcowart](https://discuss.elastic.co/u/rcowart)\
**Post date:** [March 31, 2021, 7:49am UTC](https://discuss.elastic.co/t/logstash-kafka-input-conditonnal-consuming/268664/8 "2021-03-31T07:49:51Z")

</div>

You haven't really filtered on anything. At the moment it doesn't look like you are thinking about this problem the right way. I believe what you really are trying to build is this...**[kafka\_logstash\_siem.pdf](https://github.com/elastiflow/elastiflow_for_elasticsearch/files/6234707/kafka_logstash_siem.pdf)**

---

<div class="post-metadata">

**Author:** ![Travis](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/travis/32/54079_2.png) [@Travis](https://discuss.elastic.co/u/Travis)\
**Post date:** [March 31, 2021, 2:40pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-conditonnal-consuming/268664/9 "2021-03-31T14:40:13Z")

</div>

Not agree with that : I filtered using the tag :

```
output {
    if "Fortigate" in [tags] {
        elasticsearch {

```

If I put an invalid tag, no data is written to Elasticsearch so for me it's filtering/working

Your project looks great. Will take some time to have a look on it !

---

<div class="post-metadata">

**Author:** ![rcowart](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/rcowart/32/88091_2.png) [@rcowart](https://discuss.elastic.co/u/rcowart)\
**Post date:** [March 31, 2021, 2:54pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-conditonnal-consuming/268664/10 "2021-03-31T14:54:22Z")

</div>

Not really. In your kafka input you assign the tag "Fortigate" to all of the messages consumed from the "firewall" topic. Then in the output you check to see if "Fortigate" is a tag. Of course it is _always_ a tag because you assigned it to every event in the input. So the end result is that you haven't filtered anything.

---

<div class="post-metadata">

**Author:** ![Travis](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/travis/32/54079_2.png) [@Travis](https://discuss.elastic.co/u/Travis)\
**Post date:** [March 31, 2021, 3:03pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-conditonnal-consuming/268664/11 "2021-03-31T15:03:55Z")

</div>

Yes... You've got a point ! 🙂 I didn't updated my post but after some tests I removed the tag at this input and I added it on the logstash which as a collector :

Fortigate -\> Logstash collector (where I add the tag for the fortigate syslog input) -\> Kakfa -\> Logstash (which do processing and use the tag added in the previous logstash for the kafka input)

---

<div class="post-metadata">

**Author:** ![rcowart](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/rcowart/32/88091_2.png) [@rcowart](https://discuss.elastic.co/u/rcowart)\
**Post date:** [March 31, 2021, 3:17pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-conditonnal-consuming/268664/12 "2021-03-31T15:17:52Z")

</div>

If the collector is adding the tag, why use a tag at all? Why not just produce the record to a Kafka topic called "fortigate"? The other pipeline consumes from the "fortigate" topic. You then no longer need a filter in the output, because this pipeline gets ONLY fortigate events. You also avoid consuming other logs that aren't fortigate and having to discard them.

---

<div class="post-metadata">

**Author:** ![Travis](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/travis/32/54079_2.png) [@Travis](https://discuss.elastic.co/u/Travis)\
**Post date:** [March 31, 2021, 3:27pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-conditonnal-consuming/268664/13 "2021-03-31T15:27:08Z")

</div>

Yes. I could do that. I thought that, maybe, it was better to minimize the number of topics... To facilitate maintenance. Maybe this is not a good idea after all 🤔

---

<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:** [April 28, 2021, 3:27pm UTC](https://discuss.elastic.co/t/logstash-kafka-input-conditonnal-consuming/268664/14 "2021-04-28T15:27:14Z")

</div>

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