# Unable to push messages to Kafka streams

**URL:** <https://discuss.elastic.co/t/unable-to-push-messages-to-kafka-streams/309458>\
**Category:** Logstash\
**Created:** [July 12, 2022, 8:22pm UTC](https://discuss.elastic.co/t/unable-to-push-messages-to-kafka-streams/309458 "2022-07-12T20:22:51Z")\
**Posts on this page:** 7\
**Page:** 1

<div class="post-metadata">

**Author:** ![Rao\_Nelakurti](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/rao_nelakurti/32/98729_2.png) [@Rao\_Nelakurti](https://discuss.elastic.co/u/Rao_Nelakurti)\
**Post date:** [July 12, 2022, 8:22pm UTC](https://discuss.elastic.co/t/unable-to-push-messages-to-kafka-streams/309458/1 "2022-07-12T20:22:51Z")

</div>

Hi Team,

I'm trying to send application logs (filebeat -\> logstash-\> kafka topic) to kafka stream. For some reason I'm not able to see logs in kafka.

Here is my filebeat.conf

```auto
filebeat.spool_size: 2048
filebeat.idle_timeout: 5s
fields:
  envName: gxpbt
output.logstash:
  hosts: ['xyz.logstash.com:6045']

filebeat.inputs:
- type: log
  paths:
   - /var/log/osquery/xyz.results.log
  multiline.type: pattern
  multiline.pattern: '^{'
  multiline.negate: true
  multiline.match: after
  fields:
    host_ip: "10.0.0.1"
    hostname: "test-compute"
    logtype: osquerylog
  json.message_key: log

```

Logstash conf file:

```auto
input {
  beats {
    port => 6045
  }
}
filter {
  if [fields][logtype] == "osquerylog" {
       grok {
        match => ["[log][file][path]", "%{GREEDYDATA:message}"]
        add_tag => "osquerylog"
   }
  }
}

output {
  if "osquerylog" in [tags] {
    kafka {
    codec => "json"
    bootstrap_servers => "https://xyz.kafka.com:9092"
    ssl_truststore_location => "/etc/certs/kafkatruststore.jks"
    ssl_truststore_password => "1234567"
    ssl_truststore_type => "JKS"
    ssl_keystore_location => "/etc/certs/kafkakeystore.jks"
    ssl_keystore_password => "1234567"
    ssl_keystore_type => "JKS"
    sasl_mechanism => "PLAIN"
    security_protocol => "SASL_SSL"
    request_timeout_ms => "5000"
    sasl_jaas_config => "org.apache.kafka.common.security.plain.PlainLoginModule required username='test' password='test123';"
    topic_id => "osquery_log"
    }
  }
}

```

I'm seeing following log messages in logstash logs,

```auto
[2022-07-12T20:19:48,773][INFO][org.apache.kafka.common.security.authenticator.AbstractLogin][main] Successfully logged in.
[2022-07-12T20:19:48,846][INFO][org.apache.kafka.common.utils.AppInfoParser][main] Kafka version: 2.5.1
[2022-07-12T20:19:48,846][INFO][org.apache.kafka.common.utils.AppInfoParser][main] Kafka commitId: 0efa8fb0f4c73d92
[2022-07-12T20:19:48,846][INFO][org.apache.kafka.common.utils.AppInfoParser][main] Kafka startTimeMs: 1657657188843

```

Please help me understand why I am not able to see events in kafka.

---

<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 12, 2022, 8:44pm UTC](https://discuss.elastic.co/t/unable-to-push-messages-to-kafka-streams/309458/2 "2022-07-12T20:44:14Z")

</div>

The first step would be to add

```
output { stdout { codec => rubydebug } }

```

(or a file output) to make sure that the events are reaching logstash and have the fields that you expect.

---

<div class="post-metadata">

**Author:** ![Rao\_Nelakurti](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/rao_nelakurti/32/98729_2.png) [@Rao\_Nelakurti](https://discuss.elastic.co/u/Rao_Nelakurti)\
**Post date:** [July 15, 2022, 4:13pm UTC](https://discuss.elastic.co/t/unable-to-push-messages-to-kafka-streams/309458/3 "2022-07-15T16:13:31Z")

</div>

Thanks for your reply. I can see logs being sent, because of the kafka stream offset all events were not reaching.

I'have following filebeat.conf file,

```auto
fields:
  envName: gxpbt
output.logstash:
  hosts: ['xyz.logstash.com:6045']

filebeat.inputs:
- type: log
  paths:
   - /var/log/osquery/osquery.results.log
  multiline.type: pattern
  multiline.pattern: '^{'
  multiline.negate: true
  multiline.match: after
  fields:
    host_ip: "10.0.0.1"
    hostname: "test-compute"
    regio: "ashburn"  
  json.message_key: log'
- type: log
  paths:
   - /var/log/hostname/forwarded-logs.log
  fields:
    logtype: rsyslog
- type: log
  paths:
    - /u01/data/domains/*_domain/servers/*/logs/access.log
  fields:
    logtype: wlsAccess

```

Logstash.conf:

```auto
input {
  beats {
    port => 6045
  }
}

  if [fields][region] == "ashburn" {
       grok {
        match => ["message", "%{GREEDYDATA:message}"]
        add_tag => "osquerylog"
   }
  }
 
  if [fields][logtype] == "wlsAccess" {
        grok {
                match => ["[log][file][path]","%{GREEDYDATA}/%{DATA:servername}/logs/%{DATA:dirname}/%{GREEDYDATA:filename}"]
                }
  }
  if [fields][logtype] == "rsyslog" {
       grok {
        match => ["message", "%{GREEDYDATA:message}"]
        add_tag => "rsyslog"
   }
  }

  else {
   grok {
    match => ["[log][file][path]","%{GREEDYDATA}/%{DATA:servername}/logs/%{GREEDYDATA:filename}"]
   }
  }

output {
  if "osquerylog" in [tags] {
    kafka {
    codec => "json"
    bootstrap_servers => "https://xyz.kafka.com:9092"
    ssl_truststore_location => "/etc/certs/kafkatruststore.jks"
    ssl_truststore_password => "1234567"
    ssl_truststore_type => "JKS"
    ssl_keystore_location => "/etc/certs/kafkakeystore.jks"
    ssl_keystore_password => "1234567"
    ssl_keystore_type => "JKS"
    sasl_mechanism => "PLAIN"
    security_protocol => "SASL_SSL"
    request_timeout_ms => "5000"
    sasl_jaas_config => "org.apache.kafka.common.security.plain.PlainLoginModule required username='test' password='test123';"
    topic_id => "osquery_log"
    }
  }
  if "rsyslog" in [tags] {
       syslog {
        appname => "SIEM"
        host => "10.x.x.52"
        port => "514"
        protocol => "tcp"
        codec => line { format => "%{message}" }
    }
  }
  stdout{}
  else{
   file {
     codec => line {
       format => "%{[message]}"
     }
     path => "/mnt/shared_fs_03/%{[fields][envName]}/%{[host][name]}/%{[servername]}/%{[filename]}-%{+YYYY-MM-dd}"
   }
    #stdout{}
 }
}

```

Could you please confirm if I can use **IF** and **else** conditions and filters?

---

<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 15, 2022, 4:23pm UTC](https://discuss.elastic.co/t/unable-to-push-messages-to-kafka-streams/309458/4 "2022-07-15T16:23:15Z")

</div>

> [@Rao\_Nelakurti](#):
>
> Could you please confirm if I can use **IF** and **else** conditions and filters?

Yes, logstash supports [conditionals](https://www.elastic.co/guide/en/logstash/current/event-dependent-configuration.html#conditionals) in the filter section.

---

<div class="post-metadata">

**Author:** ![Rao\_Nelakurti](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/rao_nelakurti/32/98729_2.png) [@Rao\_Nelakurti](https://discuss.elastic.co/u/Rao_Nelakurti)\
**Post date:** [July 15, 2022, 4:35pm UTC](https://discuss.elastic.co/t/unable-to-push-messages-to-kafka-streams/309458/5 "2022-07-15T16:35:50Z")

</div>

@Badger  
For example I have field "region is ashburn" under filebeat.conf. can I use the if condition as above. Can I write if else conditions as above?

---

<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 15, 2022, 4:39pm UTC](https://discuss.elastic.co/t/unable-to-push-messages-to-kafka-streams/309458/6 "2022-07-15T16:39:45Z")

</div>

Provided that the names match (region vs. regio) then the answer to both questions is yes.

---

<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 12, 2022, 4:39pm UTC](https://discuss.elastic.co/t/unable-to-push-messages-to-kafka-streams/309458/7 "2022-08-12T16:39:49Z")

</div>

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