Building a change data capture service with Debezium and Spring Boot

Modern applications often need database changes to reach search indexes, reporting tools, caches, and other services without repeatedly querying entire tables. Change data capture (CDC) solves this by recording inserts, updates, and deletes as an event stream.

Debezium reads database transaction logs and publishes structured events to Apache Kafka. A Spring Boot application can then consume those events, apply business rules, and update downstream systems. This approach is useful for Java teams building reliable integrations for customers in Sydney, Melbourne, Brisbane, and across Australia.

The design also suits Australian systems that must handle intermittent connectivity, high-volume retail activity, and privacy obligations under the Privacy Act 1988. Keeping the database as the source of truth while distributing events helps services scale without placing unnecessary load on production tables.

How Debezium fits into the architecture

A typical deployment contains a source database, Kafka, a Kafka Connect worker, a Debezium connector, and one or more Spring Boot consumers. The connector watches the database transaction log rather than executing frequent polling queries. Each committed change becomes an event in a Kafka topic.

For PostgreSQL, Debezium reads the write-ahead log. For MySQL, it reads the binary log. The exact setup differs by database, but the principle remains the same: committed changes are captured with metadata such as the table name, operation type, transaction identifier, and before-and-after values.

Kafka Connect runs Debezium outside the application process, which keeps capture responsibilities separate from business processing. Spring Boot should focus on consuming events, validating their payloads, and coordinating the required application actions.

Preparing the database and local environment

Start with a database user that has only the permissions required for replication. PostgreSQL needs logical replication enabled and a suitable replication slot. MySQL requires row-based binary logging and a user with replication privileges. The database configuration should be tested before the connector is deployed.

For local development, Docker Compose is a practical choice. Run PostgreSQL or MySQL, Kafka, Kafka Connect, and a small Spring Boot service together. A developer in Melbourne can reproduce the same stack as a teammate in Perth, reducing environment-specific surprises during integration work.

A connector configuration commonly includes the database host, credentials, server name, plugin, and the list of tables to capture. Restricting table.include.list is valuable because it avoids publishing internal tables, audit noise, or personal information that downstream consumers do not need.

Creating the Debezium connector

The connector is usually registered through Kafka Connect’s REST API. A simplified PostgreSQL configuration looks like this:

{
  "name": "orders-connector",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "postgres",
    "database.port": "5432",
    "database.user": "cdc_user",
    "database.password": "secret",
    "database.dbname": "shop",
    "topic.prefix": "shop",
    "plugin.name": "pgoutput",
    "table.include.list": "public.orders"
  }
}

When the connector starts for the first time, it can take an initial snapshot before streaming new changes. This is useful when a consumer needs an existing catalogue, although large tables may require a carefully planned snapshot window. After the snapshot, Debezium continues from the database log position.

The resulting topic may be named shop.public.orders. Its key often contains the primary key, while the value includes before, after, source, and op fields. The operation code commonly identifies create, update, delete, and read events.

Consuming events with Spring Boot

Spring Kafka provides a straightforward consumer implementation. Configure the broker address, consumer group, JSON deserialisation, and error handling in application.yml. A listener can receive the event as a JSON tree when the schema varies between tables.

@KafkaListener(topics = "shop.public.orders",
               groupId = "order-indexer")
public void receive(ConsumerRecord<String, JsonNode> record) {
    JsonNode event = record.value();
    String operation = event.path("op").asText();

    if ("c".equals(operation) || "u".equals(operation)) {
        orderIndex.save(event.path("after"));
    } else if ("d".equals(operation)) {
        orderIndex.delete(event.path("before"));
    }
}

In production, use explicit schemas or strongly typed event classes where possible. A useful code snippet guide can help keep examples and reusable Java fragments organised as the service grows.

Handling delivery, ordering, and failures

CDC consumers should be designed for at-least-once delivery. Kafka may deliver a message again after a timeout or application restart, so operations must be idempotent. An index update can safely overwrite the same document, while a payment or email action needs a deduplication key based on the event ID or transaction metadata.

Ordering is generally preserved within a Kafka partition, not across every topic. Use a stable record key, usually the database primary key, so changes for one entity remain ordered. Avoid assuming that an update received today represents the latest state unless the consumer checks version fields or source timestamps.

Failed messages should be retried with sensible back-off and eventually routed to a dead-letter topic. Monitoring should include consumer lag, connector status, replication slot growth, processing latency, and the age of the oldest unprocessed event. These metrics are especially important for Australian retailers handling sharp traffic increases around Boxing Day and major online sales.

Operating the service responsibly

Use TLS and authentication between Kafka, Kafka Connect, databases, and Spring Boot consumers. Do not place passwords in committed configuration files; use environment variables, a secrets manager, or a managed platform. Keep personal data out of event payloads unless the downstream purpose is clear and authorised.

The Australian Privacy Principles encourage careful handling of personal information, including access control, retention, transparency, and secure destruction. A service operating in Sydney or Adelaide may also need to document where replicated data is stored, particularly when cloud infrastructure crosses national borders.

Schema evolution deserves equal attention. Adding nullable fields is usually safer than renaming or removing fields immediately. Debezium’s schema history and Kafka topic compatibility settings should be backed up and tested during deployments so a consumer can recover without losing its position.

Practical decisions before production

A small proof of concept can validate the database connector, event shape, and consumer behaviour before the team commits to a larger platform. Test inserts, updates, deletes, rollbacks, connector restarts, database failover, and malformed messages.

Useful production decisions include:

With these controls in place, Debezium and Spring Boot provide a maintainable event-driven bridge between transactional Java applications and the systems that depend on timely database changes.