Streaming reference architecture built around Kafka.
See the codeJoin us on slack
Streaming reference architecture built around Kafka.

A collection of components to build a real time ingestion pipeline.
Please take a moment and read the documentation and make sure the software prerequisites are met!!
| Connector | Type | Description | Docs |
|---|---|---|---|
| AzureDocumentDb | Sink | Kafka connect Azure DocumentDb sink to subscribe to write to the cloud Azure Document Db. | Docs |
| BlockChain | Source | Kafka connect Blockchain source to subscribe to Blockchain streams and write to Kafka. | Docs |
| Bloomberg | Source | Kafka connect source to subscribe to Bloomberg streams and write to Kafka. | Docs |
| Cassandra | Source | Kafka connect Cassandra source to read Cassandra and write to Kafka. | Docs |
| Coap | Source | Kafka connect Coap source to read from IoT Coap endpoints using Californium. | Docs |
| Coap | Sink | Kafka connect Coap sink to write kafka topic payload to IoT Coap endpoints using Californium. | Docs |
| *DSE Cassandra | Sink | Certified DSE Kafka connect Cassandra sink task to write Kafka topic payloads to Cassandra. | Docs |
| Druid | Sink | Kafka connect Druid sink to write Kafka topic payloads to Druid. | Docs |
| Elastic | Sink | Kafka connect Elastic Search sink to write Kafka topic payloads to Elastic Search. | Docs |
| FTP/HTTP | Source | Kafka connect FTP and HTTP source to write file data into Kafka topics. | Docs |
| HBase | Sink | Kafka connect HBase sink to write Kafka topic payloads to HBase. | Docs |
| Hazelcast | Sink | Kafka connect Hazelcast sink to write Kafka topic payloads to Hazelcast. | Docs |
| Kudu | Sink | Kafka connect Kudu sink to write Kafka topic payloads to Kudu. | Docs |
| InfluxDb | Sink | Kafka connect InfluxDb sink to write Kafka topic payloads to InfluxDb. | Docs |
| JMS | Source | Kafka connect JMS source to write from JMS to Kafka topics. | Docs |
| JMS | Sink | Kafka connect JMS sink to write Kafka topic payloads to JMS. | Docs |
| MongoDB | Sink | Kafka connect MongoDB sink to write Kafka topic payloads to MongoDB. | Docs |
| MQTT | Source | Kafka connect MQTT source to write data from MQTT to Kafka. | Docs |
| MQTT | Sink | Kafka connect MQTT sink to write data from Kafka to MQTT. | Docs |
| Redis | Sink | Kafka connect Redis sink to write Kafka topic payloads to Redis. | Docs |
| ReThinkDB | Source | Kafka connect RethinkDb source subscribe to ReThinkDB changefeeds and write to Kafka. | Docs |
| ReThinkDB | Sink | Kafka connect RethinkDb sink to write Kafka topic payloads to RethinkDb. | Docs |
| Yahoo Finance | Source | Kafka connect Yahoo Finance source to write to Kafka. | Docs |
| VoltDB | Sink | Kafka connect Voltdb sink to write Kafka topic payloads to Voltdb. | Docs |
3.0.1 (Pending)
ftp.protocol introduced, either ftp (default) or ftps.3.0.0
0.2.6
connect.progress.enabled which will periodically report log messages processedconnect.documentdb.db to connect.documentdb.dbconnect.documentdb.database.create to connect.documentdb.db.createconnect.cassandra.source.kcql to connect.cassandra.kcqlconnect.cassandra.source.timestamp.type to connect.cassandra.timestamp.typeconnect.cassandra.source.import.poll.interval to connect.cassandra.import.poll.intervalconnect.cassandra.source.error.policy to connect.cassandra.error.policyconnect.cassandra.source.max.retries to connect.cassandra.max.retriesconnect.cassandra.source.retry.interval to connect.cassandra.retry.intervalconnect.cassandra.sink.kcql to connect.cassandra.kcqlconnect.cassandra.sink.error.policy to connect.cassandra.error.policyconnect.cassandra.sink.max.retries to connect.cassandra.max.retriesconnect.cassandra.sink.retry.interval to connect.cassandra.retry.intervalconnect.coap.bind.port to connect.coap.portconnect.coap.bind.port to connect.coap.portconnect.coap.bind.host to connect.coap.hostconnect.coap.bind.host to connect.coap.hostconnect.mongo.database to connect.mongo.dbconnect.mongo.sink.batch.size to connect.mongo.batch.sizeconnect.druid.sink.kcql to connect.druid.kcqlconnect.druid.sink.conf.file to connect.druid.kcqlconnect.druid.sink.write.timeout to connect.druid.write.timeoutconnect.elastic.sink.kcql to connect.elastic.kcqlconnect.hbase.sink.column.family to connect.hbase.column.familyconnect.hbase.sink.kcql to connect.hbase.kcqlconnect.hbase.sink.error.policy to connect.hbase.error.policyconnect.hbase.sink.max.retries to connect.hbase.max.retriesconnect.hbase.sink.retry.interval to connect.hbase.retry.intervalconnect.influx.sink.kcql to connect.influx.kcqlconnect.influx.connection.user to connect.influx.usernameconnect.influx.connection.password to connect.influx.passwordconnect.influx.connection.database to connect.influx.dbconnect.influx.connection.url to connect.influx.urlconnect.kudu.sink.kcql to connect.kudu.kcqlconnect.kudu.sink.error.policy to connect.kudu.error.policyconnect.kudu.sink.retry.interval to connect.kudu.retry.intervalconnect.kudu.sink.max.retries to connect.kudu.max.retiesconnect.kudu.sink.schema.registry.url to connect.kudu.schema.registry.urlconnect.redis.connection.password to connect.redis.passwordconnect.redis.sink.kcql to connect.redis.kcqlconnect.redis.connection.host to connect.redis.hostconnect.redis.connection.port to connect.redis.portconnect.rethink.source.host to connect.rethink.hostconnect.rethink.source.port to connect.rethink.portconnect.rethink.source.db to connect.rethink.dbconnect.rethink.source.kcql to connect.rethink.kcqlconnect.rethink.sink.host to connect.rethink.hostconnect.rethink.sink.port to connect.rethink.portconnect.rethink.sink.db to connect.rethink.dbconnect.rethink.sink.kcql to connect.rethink.kcqlconnect.jms.user to connect.jms.usernameconnect.jms.source.converters to connect.jms.convertersconnect.jms.converters and replace my kcql withConvertersconnect.jms.queues and replace my kcql withType QUEUEconnect.jms.topics and replace my kcql withType TOPICconnect.mqtt.source.kcql to connect.mqtt.kcqlconnect.mqtt.user to connect.mqtt.usernameconnect.mqtt.hosts to connect.mqtt.connection.hostsconnect.mqtt.converters and replace my kcql withConvertersconnect.mqtt.queues and replace my kcql withType=QUEUEconnect.mqtt.topics and replace my kcql withType=TOPICconnect.hazelcast.sink.kcql to connect.hazelcast.kcqlconnect.hazelcast.sink.group.name to connect.hazelcast.group.nameconnect.hazelcast.sink.group.password to connect.hazelcast.group.passwordconnect.hazelcast.sink.cluster.members tp connect.hazelcast.cluster.membersconnect.hazelcast.sink.batch.size to connect.hazelcast.batch.sizeconnect.hazelcast.sink.error.policy to connect.hazelcast.error.policyconnect.hazelcast.sink.max.retries to connect.hazelcast.max.retriesconnect.hazelcast.sink.retry.interval to connect.hazelcast.retry.intervalconnect.volt.sink.kcql to connect.volt.kcqlconnect.volt.sink.connection.servers to connect.volt.serversconnect.volt.sink.connection.user to connect.volt.usernameconnect.volt.sink.connection.password to connect.volt.passwordconnect.volt.sink.error.policy to connect.volt.error.policyconnect.volt.sink.max.retries to connect.volt.max.retriesconnect.volt.sink.retry.interval to connect.volt.retry.interval0.2.5 (8 Apr 2017)
withunwraptimestamp in the Cassandra Source for timestamp tracking.0.2.4 (26 Jan 2017)
SELECT * FROM influx-topic WITHTIMESTAMP sys_time() WITHTAG(field1, CONSTANT_KEY1=CONSTANT_VALUE1, field2,CONSTANT_KEY2=CONSTANT_VALUE1)ALL. Use connect.influx.consistency.level to set it to ONE/QUORUM/ALL/ANYconnect.influx.sink.route.query was renamed to connect.influx.sink.kcql0.2.3 (5 Jan 2017)
Struct, Schema.STRING and Json with schema in the Cassandra, ReThinkDB, InfluxDB and MongoDB sinks.export.query.route to sink.kcql.import.query.route to source.kcql.STOREAS so specify target sink types, e.g. Redis Sorted Sets, Hazelcast map, queues, ringbuffers.Requires gradle 3.0 to build.
To build
gradle compile
To test
gradle test
To create a fat jar
gradle shadowJar
You can also use the gradle wrapper
./gradlew shadowJar
To view dependency trees
gradle dependencies # or
gradle :kafka-connect-cassandra:dependencies
To build a particular project
gradle :kafka-connect-elastic5:build
We'd love to accept your contributions! Please use GitHub pull requests: fork the repo, develop and test your code, semantically commit and submit a pull request. Thanks!
Scala
99.3%
Streaming reference architecture built around Kafka.
See the codeJoin us on slack
Streaming reference architecture built around Kafka.

A collection of components to build a real time ingestion pipeline.
Please take a moment and read the documentation and make sure the software prerequisites are met!!
| Connector | Type | Description | Docs |
|---|---|---|---|
| AzureDocumentDb | Sink | Kafka connect Azure DocumentDb sink to subscribe to write to the cloud Azure Document Db. | Docs |
| BlockChain | Source | Kafka connect Blockchain source to subscribe to Blockchain streams and write to Kafka. | Docs |
| Bloomberg | Source | Kafka connect source to subscribe to Bloomberg streams and write to Kafka. | Docs |
| Cassandra | Source | Kafka connect Cassandra source to read Cassandra and write to Kafka. | Docs |
| Coap | Source | Kafka connect Coap source to read from IoT Coap endpoints using Californium. | Docs |
| Coap | Sink | Kafka connect Coap sink to write kafka topic payload to IoT Coap endpoints using Californium. | Docs |
| *DSE Cassandra | Sink | Certified DSE Kafka connect Cassandra sink task to write Kafka topic payloads to Cassandra. | Docs |
| Druid | Sink | Kafka connect Druid sink to write Kafka topic payloads to Druid. | Docs |
| Elastic | Sink | Kafka connect Elastic Search sink to write Kafka topic payloads to Elastic Search. | Docs |
| FTP/HTTP | Source | Kafka connect FTP and HTTP source to write file data into Kafka topics. | Docs |
| HBase | Sink | Kafka connect HBase sink to write Kafka topic payloads to HBase. | Docs |
| Hazelcast | Sink | Kafka connect Hazelcast sink to write Kafka topic payloads to Hazelcast. | Docs |
| Kudu | Sink | Kafka connect Kudu sink to write Kafka topic payloads to Kudu. | Docs |
| InfluxDb | Sink | Kafka connect InfluxDb sink to write Kafka topic payloads to InfluxDb. | Docs |
| JMS | Source | Kafka connect JMS source to write from JMS to Kafka topics. | Docs |
| JMS | Sink | Kafka connect JMS sink to write Kafka topic payloads to JMS. | Docs |
| MongoDB | Sink | Kafka connect MongoDB sink to write Kafka topic payloads to MongoDB. | Docs |
| MQTT | Source | Kafka connect MQTT source to write data from MQTT to Kafka. | Docs |
| MQTT | Sink | Kafka connect MQTT sink to write data from Kafka to MQTT. | Docs |
| Redis | Sink | Kafka connect Redis sink to write Kafka topic payloads to Redis. | Docs |
| ReThinkDB | Source | Kafka connect RethinkDb source subscribe to ReThinkDB changefeeds and write to Kafka. | Docs |
| ReThinkDB | Sink | Kafka connect RethinkDb sink to write Kafka topic payloads to RethinkDb. | Docs |
| Yahoo Finance | Source | Kafka connect Yahoo Finance source to write to Kafka. | Docs |
| VoltDB | Sink | Kafka connect Voltdb sink to write Kafka topic payloads to Voltdb. | Docs |
3.0.1 (Pending)
ftp.protocol introduced, either ftp (default) or ftps.3.0.0
0.2.6
connect.progress.enabled which will periodically report log messages processedconnect.documentdb.db to connect.documentdb.dbconnect.documentdb.database.create to connect.documentdb.db.createconnect.cassandra.source.kcql to connect.cassandra.kcqlconnect.cassandra.source.timestamp.type to connect.cassandra.timestamp.typeconnect.cassandra.source.import.poll.interval to connect.cassandra.import.poll.intervalconnect.cassandra.source.error.policy to connect.cassandra.error.policyconnect.cassandra.source.max.retries to connect.cassandra.max.retriesconnect.cassandra.source.retry.interval to connect.cassandra.retry.intervalconnect.cassandra.sink.kcql to connect.cassandra.kcqlconnect.cassandra.sink.error.policy to connect.cassandra.error.policyconnect.cassandra.sink.max.retries to connect.cassandra.max.retriesconnect.cassandra.sink.retry.interval to connect.cassandra.retry.intervalconnect.coap.bind.port to connect.coap.portconnect.coap.bind.port to connect.coap.portconnect.coap.bind.host to connect.coap.hostconnect.coap.bind.host to connect.coap.hostconnect.mongo.database to connect.mongo.dbconnect.mongo.sink.batch.size to connect.mongo.batch.sizeconnect.druid.sink.kcql to connect.druid.kcqlconnect.druid.sink.conf.file to connect.druid.kcqlconnect.druid.sink.write.timeout to connect.druid.write.timeoutconnect.elastic.sink.kcql to connect.elastic.kcqlconnect.hbase.sink.column.family to connect.hbase.column.familyconnect.hbase.sink.kcql to connect.hbase.kcqlconnect.hbase.sink.error.policy to connect.hbase.error.policyconnect.hbase.sink.max.retries to connect.hbase.max.retriesconnect.hbase.sink.retry.interval to connect.hbase.retry.intervalconnect.influx.sink.kcql to connect.influx.kcqlconnect.influx.connection.user to connect.influx.usernameconnect.influx.connection.password to connect.influx.passwordconnect.influx.connection.database to connect.influx.dbconnect.influx.connection.url to connect.influx.urlconnect.kudu.sink.kcql to connect.kudu.kcqlconnect.kudu.sink.error.policy to connect.kudu.error.policyconnect.kudu.sink.retry.interval to connect.kudu.retry.intervalconnect.kudu.sink.max.retries to connect.kudu.max.retiesconnect.kudu.sink.schema.registry.url to connect.kudu.schema.registry.urlconnect.redis.connection.password to connect.redis.passwordconnect.redis.sink.kcql to connect.redis.kcqlconnect.redis.connection.host to connect.redis.hostconnect.redis.connection.port to connect.redis.portconnect.rethink.source.host to connect.rethink.hostconnect.rethink.source.port to connect.rethink.portconnect.rethink.source.db to connect.rethink.dbconnect.rethink.source.kcql to connect.rethink.kcqlconnect.rethink.sink.host to connect.rethink.hostconnect.rethink.sink.port to connect.rethink.portconnect.rethink.sink.db to connect.rethink.dbconnect.rethink.sink.kcql to connect.rethink.kcqlconnect.jms.user to connect.jms.usernameconnect.jms.source.converters to connect.jms.convertersconnect.jms.converters and replace my kcql withConvertersconnect.jms.queues and replace my kcql withType QUEUEconnect.jms.topics and replace my kcql withType TOPICconnect.mqtt.source.kcql to connect.mqtt.kcqlconnect.mqtt.user to connect.mqtt.usernameconnect.mqtt.hosts to connect.mqtt.connection.hostsconnect.mqtt.converters and replace my kcql withConvertersconnect.mqtt.queues and replace my kcql withType=QUEUEconnect.mqtt.topics and replace my kcql withType=TOPICconnect.hazelcast.sink.kcql to connect.hazelcast.kcqlconnect.hazelcast.sink.group.name to connect.hazelcast.group.nameconnect.hazelcast.sink.group.password to connect.hazelcast.group.passwordconnect.hazelcast.sink.cluster.members tp connect.hazelcast.cluster.membersconnect.hazelcast.sink.batch.size to connect.hazelcast.batch.sizeconnect.hazelcast.sink.error.policy to connect.hazelcast.error.policyconnect.hazelcast.sink.max.retries to connect.hazelcast.max.retriesconnect.hazelcast.sink.retry.interval to connect.hazelcast.retry.intervalconnect.volt.sink.kcql to connect.volt.kcqlconnect.volt.sink.connection.servers to connect.volt.serversconnect.volt.sink.connection.user to connect.volt.usernameconnect.volt.sink.connection.password to connect.volt.passwordconnect.volt.sink.error.policy to connect.volt.error.policyconnect.volt.sink.max.retries to connect.volt.max.retriesconnect.volt.sink.retry.interval to connect.volt.retry.interval0.2.5 (8 Apr 2017)
withunwraptimestamp in the Cassandra Source for timestamp tracking.0.2.4 (26 Jan 2017)
SELECT * FROM influx-topic WITHTIMESTAMP sys_time() WITHTAG(field1, CONSTANT_KEY1=CONSTANT_VALUE1, field2,CONSTANT_KEY2=CONSTANT_VALUE1)ALL. Use connect.influx.consistency.level to set it to ONE/QUORUM/ALL/ANYconnect.influx.sink.route.query was renamed to connect.influx.sink.kcql0.2.3 (5 Jan 2017)
Struct, Schema.STRING and Json with schema in the Cassandra, ReThinkDB, InfluxDB and MongoDB sinks.export.query.route to sink.kcql.import.query.route to source.kcql.STOREAS so specify target sink types, e.g. Redis Sorted Sets, Hazelcast map, queues, ringbuffers.Requires gradle 3.0 to build.
To build
gradle compile
To test
gradle test
To create a fat jar
gradle shadowJar
You can also use the gradle wrapper
./gradlew shadowJar
To view dependency trees
gradle dependencies # or
gradle :kafka-connect-cassandra:dependencies
To build a particular project
gradle :kafka-connect-elastic5:build
We'd love to accept your contributions! Please use GitHub pull requests: fork the repo, develop and test your code, semantically commit and submit a pull request. Thanks!
Scala
99.3%