# How to connect Kafka with Elasticsearch?

**URL:** https://discuss.elastic.co/t/how-to-connect-kafka-with-elasticsearch/118002
**Category:** Elasticsearch
**Created:** [February 1, 2018, 10:56am UTC](https://discuss.elastic.co/t/how-to-connect-kafka-with-elasticsearch/118002 "2018-02-01T10:56:46Z")
**Posts on this page:** 3
**Page:** 1

<div class="post-metadata">

### Author: ![f26227279](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/f26227279/32/21296_2.png) [@f26227279](https://discuss.elastic.co/u/f26227279)
#### Post date: [February 1, 2018, 10:56am UTC](https://discuss.elastic.co/t/how-to-connect-kafka-with-elasticsearch/118002/1 "2018-02-01T10:56:46Z")

</div>

Hi everyone, I am new in Kafka, I use kafka to collect netflow through logstash(it is ok), and I want to send the data to elasticsearch from kafka, but there are some problem.  
netflow to kafka logstash config:

```
input{
	udp{
		host => "120.127.XXX.XX"
		port => 5556
		codec => netflow
	}
}
	filter{
		
	}
output {
  kafka {
    bootstrap_servers => "localhost:9092"    
    topic_id => "test"    
  }
  stdout{codec=> rubydebug}
}

```

kafka to elasticsearch logstash:

```
input {
      kafka { }
    }
    output {
        elasticsearch {
            hosts => ["120.127.XXX.XX:9200"]
        }
    	stdout{codec=> rubydebug}
    }

```

log:  
D:\ELK\logstash-6.1.1\bin\>logstash -f kafkatoES.conf --path.data D:\ELK\logstash-6.1.1\datatest  
Sending Logstash's logs to D:/ELK/logstash-6.1.1/logs which is now configured via log4j2.properties  
[2018-02-01T18:52:59,713][INFO][logstash.modules.scaffold] Initializing module {:module\_name=\>"fb\_apache", :directory=\>"D:/ELK/logstash-6.1.1/modules/fb\_apache/configuration"}  
[2018-02-01T18:52:59,728][INFO][logstash.modules.scaffold] Initializing module {:module\_name=\>"netflow", :directory=\>"D:/ELK/logstash-6.1.1/modules/netflow/configuration"}  
[2018-02-01T18:53:00,072][WARN][logstash.config.source.multilocal] Ignoring the 'pipelines.yml' file because modules or command line options are specified  
[2018-02-01T18:53:01,070][INFO][logstash.runner] Starting Logstash {"logstash.version"=\>"6.1.1"}  
[2018-02-01T18:53:01,804][INFO][logstash.agent] Successfully started Logstash API endpoint {:port=\>9601}  
[2018-02-01T18:53:09,024][INFO][logstash.outputs.elasticsearch] Elasticsearch pool URLs updated {:changes=\>{:removed=\>[], :added=\>[[http://120.127.XX.XX:9200/](http://120.127.XX.XX:9200/)]}}  
[2018-02-01T18:53:09,040][INFO][logstash.outputs.elasticsearch] Running health check to see if an Elasticsearch connection is working {:healthcheck\_url=\>[http://120.127.XX.XX:9200/](http://120.127.XX.XX:9200/), :path=\>"/"}  
[2018-02-01T18:53:09,305][WARN][logstash.outputs.elasticsearch] Restored connection to ES instance {:url=\>"[http://120.127.XX.XX:9200/](http://120.127.XX.XX:9200/)"}  
[2018-02-01T18:53:09,383][INFO][logstash.outputs.elasticsearch] ES Output version determined {:es\_version=\>nil}  
[2018-02-01T18:53:09,383][WARN][logstash.outputs.elasticsearch] Detected a 6.x and above cluster: the `type` event field won't be used to determine the document \_type {:es\_version=\>6}  
[2018-02-01T18:53:09,415][INFO][logstash.outputs.elasticsearch] Using mapping template from {:path=\>nil}  
[2018-02-01T18:53:09,430][INFO][logstash.outputs.elasticsearch] Attempting to install template {:manage\_template=\>{"template"=\>"logstash-_", "version"=\>60001, "settings"=\>{"index.refresh\_interval"=\>"5s"}, "mappings"=\>{"default"=\>{"dynamic\_templates"=\>[{"message\_field"=\>{"path\_match"=\>"message", "match\_mapping\_type"=\>"string", "mapping"=\>{"type"=\>"text", "norms"=\>false}}}, {"string\_fields"=\>{"match"=\>"_", "match\_mapping\_type"=\>"string", "mapping"=\>{"type"=\>"text", "norms"=\>false, "fields"=\>{"keyword"=\>{"type"=\>"keyword", "ignore\_above"=\>256}}}}}], "properties"=\>{"@timestamp"=\>{"type"=\>"date"}, "@version"=\>{"type"=\>"keyword"}, "geoip"=\>{"dynamic"=\>true, "properties"=\>{"ip"=\>{"type"=\>"ip"}, "location"=\>{"type"=\>"geo\_point"}, "latitude"=\>{"type"=\>"half\_float"}, "longitude"=\>{"type"=\>"half\_float"}}}}}}}}  
[2018-02-01T18:53:09,493][INFO][logstash.outputs.elasticsearch] New Elasticsearch output {:class=\>"LogStash::Outputs::ElasticSearch", :hosts=\>["[//120.127.XXX.XX:9200](https://120.127.XXX.XX:9200)"]}  
[2018-02-01T18:53:09,524][INFO][logstash.pipeline] Starting pipeline {:pipeline\_id=\>"main", "pipeline.workers"=\>16, "pipeline.batch.size"=\>125, "pipeline.batch.delay"=\>5, "pipeline.max\_inflight"=\>2000, :thread=\>"#\<Thread:0x45e62903 run\>"}  
[2018-02-01T18:53:09,609][INFO][logstash.pipeline] Pipeline started {"[pipeline.id](http://pipeline.id)"=\>"main"}  
SLF4J: Class path contains multiple SLF4J bindings.  
SLF4J: Found binding in [jar:file:/D:/ELK/logstash-6.1.1/logstash-core/lib/org/apache/logging/log4j/log4j-slf4j-impl/2.6.2/log4j-slf4j-impl-2.6.2.jar!/org/slf4j/impl/StaticLoggerBinder.class]  
SLF4J: Found binding in [jar:file:/D:/ELK/logstash-6.1.1/vendor/bundle/jruby/2.3.0/gems/logstash-input-kafka-8.0.2/vendor/jar-dependencies/runtime-jars/log4j-slf4j-impl-2.8.2.jar!/org/slf4j/impl/StaticLoggerBinder.class]  
SLF4J: See [http://www.slf4j.org/codes.html#multiple\_bindings](http://www.slf4j.org/codes.html#multiple_bindings) for an explanation.  
SLF4J: Actual binding is of type [org.apache.logging.slf4j.Log4jLoggerFactory]  
[2018-02-01T18:53:09,789][INFO][logstash.agent] Pipelines running {:count=\>1, :pipelines=\>["main"]}  
[2018-02-01T18:53:09,852][INFO][org.apache.kafka.clients.consumer.ConsumerConfig] ConsumerConfig values:  
[auto.commit.interval.ms](http://auto.commit.interval.ms) = 5000  
auto.offset.reset = latest  
bootstrap.servers = [localhost:9092]  
check.crcs = true  
[client.id](http://client.id) = logstash-0  
[connections.max.idle.ms](http://connections.max.idle.ms) = 540000  
enable.auto.commit = true  
exclude.internal.topics = true  
metrics.recording.level = INFO  
[metrics.sample.window.ms](http://metrics.sample.window.ms) = 30000  
partition.assignment.strategy = [class org.apache.kafka.clients.consumer.RangeAssignor]  
receive.buffer.bytes = 65536  
[reconnect.backoff.max.ms](http://reconnect.backoff.max.ms) = 1000  
[reconnect.backoff.ms](http://reconnect.backoff.ms) = 50  
[request.timeout.ms](http://request.timeout.ms) = 305000  
[retry.backoff.ms](http://retry.backoff.ms) = 100  
sasl.jaas.config = null  
sasl.kerberos.kinit.cmd = /usr/bin/kinit  
ssl.keystore.password = null  
ssl.keystore.type = JKS  
ssl.protocol = TLS  
ssl.provider = null  
ssl.secure.random.implementation = null  
ssl.trustmanager.algorithm = PKIX  
ssl.truststore.location = null  
ssl.truststore.password = null  
ssl.truststore.type = JKS  
value.deserializer = class org.apache.kafka.common.serialization.StringDeserializer

[2018-02-01T18:53:09,945][INFO][org.apache.kafka.common.utils.AppInfoParser] Kafka version : 0.11.0.0  
[2018-02-01T18:53:09,945][INFO][org.apache.kafka.common.utils.AppInfoParser] Kafka commitId : cb8625948210849f  
[2018-02-01T18:53:10,149][INFO][org.apache.kafka.clients.consumer.internals.AbstractCoordinator] Discovered coordinator winoc-netflow:9092 (id: 2147483647 rack: null) for group logstash.  
[2018-02-01T18:53:10,164][INFO][org.apache.kafka.clients.consumer.internals.ConsumerCoordinator] Revoking previously assigned partitions [] for group logstash  
[2018-02-01T18:53:10,164][INFO][org.apache.kafka.clients.consumer.internals.AbstractCoordinator] (Re-)joining group logstash  
[2018-02-01T18:53:10,180][INFO][org.apache.kafka.clients.consumer.internals.AbstractCoordinator] Successfully joined group logstash with generation 6  
[2018-02-01T18:53:10,180][INFO][org.apache.kafka.clients.consumer.internals.ConsumerCoordinator] Setting newly assigned partitions [logstash-0] for group logstash

thank you in advance!

---

<div class="post-metadata">

### Author: ![rmoff](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/rmoff/32/20212_2.png) [@rmoff](https://discuss.elastic.co/u/rmoff)
#### Post date: [February 1, 2018, 6:43pm UTC](https://discuss.elastic.co/t/how-to-connect-kafka-with-elasticsearch/118002/2 "2018-02-01T18:43:56Z")

</div>

(I'm copying my answer from the [x-post to SO](https://stackoverflow.com/questions/48561197/how-to-connect-kafka-with-elasticsearch))

I would suggest using Kafka Connect and its [Elasticsearch sink](https://docs.confluent.io/current/connect/connect-elasticsearch/docs/). I actually presented on exactly this subject last night 🙂 [Here are the slides](https://speakerdeck.com/rmoff/building-streaming-data-pipelines-with-elasticsearch-apache-kafka-and-ksql).

You can see a detailed example [here](https://www.confluent.io/blog/blogthe-simplest-useful-kafka-connect-data-pipeline-in-the-world-or-thereabouts-part-2/).

---

<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: [March 1, 2018, 6:44pm UTC](https://discuss.elastic.co/t/how-to-connect-kafka-with-elasticsearch/118002/3 "2018-03-01T18:44:13Z")

</div>

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