CDC with Kafka Connect and Debezium
First written August 2022, last updated September 2026.
Setting up change data capture with databases, Apache Kafka, Kafka Connect, and Debezium takes time, with tricky configuration along the way. Here we'll walk through a setup of all components to show what is possible and give you the pieces to bring this into your project.
Real-Time Toolkit
This demo is available in the rdbms-cdc-nosql folder of the dev-local-demos project. It uses applications available through containers in dev-local. Within a few minutes, you can see change data capture from relational databases (Postgres, MySQL v8, and MySQL v5) into NoSQL data stores (Mongo, Cassandra, and Elasticsearch). Stream enrichment processing is done with ksqlDB.
Challenges
The specific touch-points covered here:
- Enabling logging within a database
- Connector setting nuances
- Management of Kafka Connect Secrets
- Logical Data-types
Database logging
Each database has its own nuances to set up logging. This is critical for any change data capture process, and something that needs to be well understood for success. Here are the settings and issues in the configuration of Postgres and MySQL with Debezium for change data capture. This isn't a complete overview of all the settings, but rather insight into the complexities your database operations and development teams need to work through together. Work with your database administrators to enable database logging on the tables that are needed, along with any snapshotting or specific configurations.
Postgres and Debezium
TL;DR
wal_levelmust be set tological.- If you are running Postgres with docker compose, override the command with the following:
command: "postgres -c wal_level=logical"- The connector needs to enable the
pgoutputplugin. - Add the following to the connector configuration:
"plugin.name" : "pgoutput"Details
Postgres needs to be able to capture changes; this is done through the write-ahead log (WAL).
The amount of data captured is based on the wal_level settings. The default setting is replica, but that is an insufficient level of data for Debezium. The logical setting includes replica information and additional logical change sets.
Debezium has to be configured to use the pgoutput plugin. Use the configuration property plugin.name to set this.
Troubleshooting
This is a set of errors seen when using Debezium Postgres Source Connector.
- Attempting to start Debezium with Postgres, without
wal_levelproperly defined.
Connector configuration is invalid and contains the following 1 error(s): Postgres server wal_level property must be "logical" but is: replica- Restarting Postgres without
-c wal_level=logicalwill result in Postgres failing to start with the following error:
FATAL: logical replication slot "debezium" exists, but wal_level < logical- Starting a connector without
pgoutputplugin enabled.
io.debezium.DebeziumException: Creation of replication slot failedIf (we should say when) you uncover an error, take the time to document it well. A déjà-vu moment when an error resurfaces isn't fun.
MySQL (v8) and Debezium
TL;DR
MySQL 8 has logging enabled by default. In production, however, you do need to verify this with operations.
Troubleshooting
You delete and recreate your MySQL database, but reuse the connector (and the state it persists in Kafka):
Caused by: io.debezium.DebeziumException:Client requested master to start replication from position > file size Error code: 1236; SQLSTATE: HY000MySQL (v5) and Debezium
TL;DR
- Binlog isn't enabled by default: set the
log-binname and configurebinlog_formattorow. - Ensure logs are retained longer than any reprocessing window.
- Debezium expects to resume where it left off; modifying the logs after Debezium has started can lead to unexpected errors.
Details
Add the following properties to your database’s mysql.cnf file.
server-id = 1log_bin = mysql-binexpire_logs_days = 99binlog_format = rowTroubleshooting
You recreate your database but use the same instance of the connector (v5 has the same error as v8):
Caused by: io.debezium.DebeziumException:Client requested master to start replication from position > file size Error code: 1236; SQLSTATE: HY000.MySQL was shut down or became unreachable:
Caused by: io.debezium.DebeziumException: Failed to read next byte from position XXXXXDebezium
Debezium is an excellent open-source change data capture product. It provides a lot of features that make it a powerful tool. We find a few working examples make it a lot easier to understand. Here are a few things worth knowing up front; they made it easier to configure connectors and quickly see the rewards of change data capture.
database.server.name- This property isn't a connection property to the database, but rather the name used to keep this connection uniquely identified. It's used in the topic name generated for the CDC process against the source database. Our suggestion: don't use the database type as the name (e.g. Postgres or MySQL). Picking a name like this could cause those maintaining the code to believe this name needs to align with the type of the database. Debezium 2.0 renamed this property to
topic.prefix. io.debezium.transforms.ExtractNewRecordState- By default, Debezium provides nested elements of before, after, and operation. For most use-cases, extracting just the
afterstate is sufficient, and Debezium provides a Single Message Transform (SMT) to do just that. Nested elements can be tricky if you aren't writing stream applications, so allowing the data to be flattened with one simple SMT is very helpful. Using this SMT makes it easier to pull data intoksqlDBfor enrichment. - Predicates (Apache Kafka 2.6)
- Prior to Apache Kafka 2.6, transformations were unconditional, making a single connector process multiple tables more difficult. By using predicates, a single connector can have different rules for extracting the key from the message.
- A common use of SMTs in Debezium connectors is to pull out the element from the value that is the primary key. This ensures that events on a given row (primary key) in the database are processed in order.
database.history- Many connectors allow for metadata related to the connector to be sourced to a different Kafka cluster. This flexibility leads to confusion, especially for developers new to Kafka Connect and to a specific connector.
- Debezium’s database history is designed this way. You need to set up the bootstrap servers, protocol, and other connection settings for the Kafka cluster that maintains this information, even if it is the same cluster. For enterprise deployments, this flexibility is critical. For proofs of concept, development, and getting something running quickly, it is a lot of duplicate configuration. Debezium 2.0 renamed these properties to
schema.history.internal.*. decimal.handling.mode- Setting this to
stringcan address sink connector issues that can't handle the decimal logical type. The demo code usesksqlDBto cast decimal to string, for downstream sinks that need the help, but this is an alternative approach.
Connector Secrets
- Apache Kafka provides a Config Provider interface that allows secrets to be stored separately from configuration. This is also available to a distributed connect cluster, and accessible from the connectors.
- Configure a file provider (
org.apache.kafka.common.config.provider.FileConfigProvider) in the settings of the distributed Connect cluster.
config.providers=file config.providers.file.class=org.apache.kafka.common.config.provider.FileConfigProviderIn addition to using these from your Kafka component configurations, they are also accessible from connectors, such as:
"username" : "${file:/etc/kafka-connect/secrets/mysql.properties:USERNAME}"- Put more than secrets in this file: store the database connection URL and any other settings that vary between deployment environments. Keeping those out of the configuration means a single artifact can be published and maintained with your source code.
"connection.url" : "${file:/etc/kafka-connect/secrets/mysql.properties:CONNECTION_URL}","connection.user" : "${file:/etc/kafka-connect/secrets/mysql.properties:CONNECTION_USER}","connection.password" : "${file:/etc/kafka-connect/secrets/mysql.properties:CONNECTION_PASSWORD}"Schemas and Data-types
When it comes to data-types, especially those considered to be logical data types in Connect API, not all connectors are the same. If you are doing change data capture, odds are you will have decimals and timestamps.
Fortunately, timestamps are stored as long epoch, which will usually translate into a database even if the logical type isn't properly handled. Decimals, however, are stored as byte arrays. If the connector doesn’t properly invoke the logical converter, it won't be properly converted for the end system. To make matters worse, connectors aren't consistent in how they handle the errors.
Specific Connector Observations
Mongo
In this demonstration, the MongoDB Sink Connector properly handles the logical type, and data is stored correctly. The SchemaRecordConverter properly handles the conversion, but as you can see, the converter has to account for and handle logical types; it isn't done within the Connect API.
Data shown in MongoDB:

Elasticsearch
If Elasticsearch sink connector creates the index in elastic (schema.ignore=false), logical types are handled properly. If the sink connector doesn’t create the index (schema.ignore=true), logical converters aren't processed, and logical-type decimals end up as an array of bytes in Elasticsearch.
Index generated by the connector. In each record the amount is a decimal value.

Index is built manually through the Elasticsearch API (the connector doesn't create it). The decimal bytes are passed as-is to the index, yielding an undesired result. Each record is the byte array value (the physical representation of the logical type).

The amounts shown in the Kibana screenshot of HUM=, G2Q=, and FIA= are the physical byte arrays converted to strings.
Cassandra
The Datastax Cassandra Sink Connector doesn't handle the decimal logical-type correctly. Worse, the conversion only produces a warning in the Connect cluster log. No data is written, and the connector keeps running. When checking the log, you can see:
WARN Error decoding/mapping Kafka record ... : Codec not found for requested operation:[DECIMAL <-> java.nio.ByteBuffer] (com.datastax.oss.kafka.sink.CassandraSinkTask)Sink connector checklist
- Be sure to test with decimals, timestamps, and dates, unless they truly aren’t in your use case.
- Don’t simplify your POC. For example, don’t let the Elasticsearch sink connector create your indexes in your POC if that isn't possible in production.
- Watch the logs and check for
WARNor evenINFO. - If you have issues with logical types, you can have Debezium use strings for decimals (
decimal.handling.mode=string), or you can use a stream processor (e.g. ksqlDB) to cast a decimal to a string.
Takeaways
- Enabling database logging can be tricky, and each database has its own way of configuring and enabling it.
- Having an end-to-end proof of concept showcasing change data capture is a great way to get developer buy-in and involvement, but plan time to work with your database operations team to get it enabled in the enterprise.
- Use Apache Kafka’s Config Provider for secrets and environment-specific differences, such as the database connection string. Checking a single configuration artifact into revision control, without environment specifics, is a great benefit.
- Validate decimals, dates, and timestamps, as not all connectors handle them correctly.
- Check out the dev-local project and the
rdbms-cdc-nosqldemo in dev-local-demos. The specifics discussed here are based on that demonstration, and scripts are there to get you observing change data capture with Debezium and Kafka within minutes.
Working on something like this?
Start a Conversation