9.6 Example: MySQL → Kafka → Elasticsearch
Three options: Confluent Hub client, download from Confluent Hub (or wherever the connector is hosted), or build from source:
6.1 Getting connectors
Three options: Confluent Hub client, download from Confluent Hub (or wherever the connector is hosted), or build from source:
git clone https://github.com/confluentinc/kafka-connect-elasticsearch
mvn install -DskipTests
# repeat for the JDBC connectorInstall into plugin.path:
mkdir /opt/connectors/jdbc
mkdir /opt/connectors/elastic
cp .../kafka-connect-jdbc-10.3.x-SNAPSHOT.jar /opt/connectors/jdbc
cp .../kafka-connect-elasticsearch-11.1.0-SNAPSHOT.jar /opt/connectors/elastic
cp .../kafka-connect-elasticsearch-11.1.0-SNAPSHOT-package/share/java/\
kafka-connect-elasticsearch/* /opt/connectors/elastic⚠️ "since we need to connect not just to any database but specifically to MySQL, you'll need to download and install a MySQL JDBC driver. THE DRIVER DOESN'T SHIP WITH THE CONNECTOR FOR LICENSE REASONS." → place the jar in
/opt/connectors/jdbc.
Restart workers and confirm:
curl http://localhost:8083/connector-plugins
# ElasticsearchSinkConnector (sink), JdbcSinkConnector (sink),
# JdbcSourceConnector (source)6.2 Source data
create database test;
use test;
create table login (username varchar(30), login_time datetime);
insert into login values ('gwenshap', now());
insert into login values ('tpalino', now());6.3 💡 Discovering configuration via the REST API
You don't have to read the docs — ask the API:
curl -X PUT -d '{"connector.class":"JdbcSource"}' \
localhost:8083/connector-plugins/JdbcSourceConnector/config/validate/ \
--header "content-Type:application/json"{ "configs": [ { "definition": {
"default_value": "",
"display_name": "Timestamp Column Name",
"documentation": "The name of the timestamp column to use to detect
new or modified rows. This column may not be nullable.",
"group": "Mode",
"importance": "MEDIUM",
"name": "timestamp.column.name",
"order": 3, "required": false, "type": "STRING", "width": "MEDIUM" } },
... ] }"We asked the REST API to validate configuration for a connector and sent it a configuration with just the class name (this is the bare minimum configuration necessary). As a response, we got the JSON definition of ALL AVAILABLE CONFIGURATIONS."
This is a genuinely useful trick: the validate endpoint doubles as self-documenting config discovery, including defaults, importance, and grouping.
6.4 The JDBC source connector
echo '{"name":"mysql-login-connector", "config":{
"connector.class":"JdbcSourceConnector",
"connection.url":"jdbc:mysql://127.0.0.1:3306/test?user=root",
"mode":"timestamp",
"table.whitelist":"login",
"validate.non.null":false,
"timestamp.column.name":"login_time",
"topic.prefix":"mysql."}}' | curl -X POST -d @- \
http://localhost:8083/connectors --header "content-Type:application/json"Verify:
bin/kafka-console-consumer.sh --bootstrap-server=localhost:9092 \
--topic mysql.login --from-beginningTroubleshooting — check the Connect worker logs:
[2016-10-16 19:39:40,482] ERROR Error while starting connector
mysql-login-connector (org.apache.kafka.connect.runtime.WorkerConnector:108)
org.apache.kafka.connect.errors.ConnectException: java.sql.SQLException:
Access denied for user 'root;'@'localhost' (using password: NO)"Other issues can involve the existence of the driver in the classpath or permissions to read the table."
Then: "if you insert additional rows in the login table, you should immediately see them reflected in the mysql.login topic."
⚠️ CHANGE DATA CAPTURE AND THE DEBEZIUM PROJECT
*"The JDBC connector we are using uses JDBC and SQL to SCAN database tables for new records. It detects new records by using timestamp fields or an incrementing primary key. THIS IS A RELATIVELY INEFFICIENT AND AT TIMES INACCURATE PROCESS.
All relational databases have a TRANSACTION LOG (also called redo log, binlog, or write-ahead log) as part of their implementation, and many allow external systems to read data directly from their transaction log — a FAR MORE ACCURATE AND EFFICIENT process known as CHANGE DATA CAPTURE. Most modern ETL systems depend on change data capture as a data source.
*"The Debezium Project provides a collection of high-quality, open source, change capture connectors for a variety of databases. If you are planning on streaming data from a relational database to Kafka, we HIGHLY RECOMMEND using a Debezium change capture connector if one exists for your database.
In addition, the Debezium documentation is one of the best we've seen — it covers useful design patterns and use cases related to change data capture, especially in the context of microservices."*
| JDBC polling | CDC (Debezium) |
|---|---|
SELECT ... WHERE ts > last_ts | reads the binlog / WAL directly |
| ✗ misses DELETEs entirely | ✓ sees inserts, updates, and deletes |
| ✗ misses rows updated within the same timestamp granularity | ✓ exact transaction ordering |
| ✗ misses intermediate values between polls | ✓ every intermediate value |
| ✗ full-table scan load on the database | ✓ near-zero query load |
► Polling is “relatively inefficient and at times inaccurate”.
6.5 The Elasticsearch sink connector
elasticsearch &
curl http://localhost:9200/ # verify it's upecho '{"name":"elastic-login-connector", "config":{
"connector.class":"ElasticsearchSinkConnector",
"connection.url":"http://localhost:9200",
"type.name":"mysql-data",
"topics":"mysql.login",
"key.ignore":true}}' | curl -X POST -d @- \
http://localhost:8083/connectors --header "content-Type:application/json"The configs explained — and key.ignore is the interesting one:
- "
connection.urlis simply the URL of the local Elasticsearch server."- "Each topic in Kafka will become, by default, a separate Elasticsearch index, with the same name as the topic."
- "The JDBC connector DOES NOT POPULATE THE MESSAGE KEY. As a result, the events in Kafka have NULL KEYS. Because the events in Kafka lack keys, we need to tell the Elasticsearch connector to use the topic name, partition ID, and offset as the key for each event. This is done by setting
key.ignoreto true."
Document _id becomes: "mysql.login+0+0"
└topic──┘ └p┘└off┘- ► This is why the search results show
_idvalues like"mysql.login+0+1". - ► It also makes the sink idempotent: reprocessing the same Kafka record overwrites the same Elasticsearch document instead of duplicating it. (This is the “exactly-once via unique keys” pattern from Ch. 8 §2.2.)
Verify the index and search it:
$ curl 'localhost:9200/_cat/indices?v'
health status index ... docs.count ... store.size
yellow open mysql.login ... 2 ... 3.9kb
$ curl -s -X "GET" "http://localhost:9200/mysql.login/_search?pretty=true"
# hits: _id "mysql.login+0+0" → {"username":"gwenshap","login_time":1621699811000}
# _id "mysql.login+0+1" → {"username":"tpalino", "login_time":1621699816000}"If the index isn't there, look for errors in the Connect worker log. Missing configurations or libraries are common causes for errors."
BUILD YOUR OWN CONNECTORS
"The Connector API is public and anyone can create a new connector. So if the datastore you wish to integrate with does not have an existing connector, we encourage you to write your own. You can then contribute it to Confluent Hub so others can discover and use it." Resources: multiple blog posts, talks from Kafka Summit NY 2019, Kafka Summit London 2018, ApacheCon, existing connectors as a starting point, and an Apache Maven archetype to jump-start. Ask on
users@kafka.apache.org.