Table of contents
Open Table of contents
- Introduction
- Background
- Solutions
- How Does CDC—Change Data Capture—Solve These Problems?
- Scenario Two Example: Oracle-to-PostgreSQL Synchronization Through CDC
- Approaches for Scenarios One and Three
- ETL (Extract, Transform, and Load) and Flink CDC
- Using CDC for Cache Consistency
- Closing Thoughts
- References
Introduction
I strongly recommend this excellent Confluent CDC article, which includes sample code. After reading it, you may not need most of what I write below.
Background
Some business scenarios require synchronization across systems. Typical examples include:
- Synchronizing a database table to Elasticsearch to improve search speed.
- Synchronizing the same table across databases, either of the same type or, as in this article, different types.
- Synchronizing different tables in one database. For example, inserting into A may also require operations on B.
- Database/cache synchronization, or cache consistency, for frequently read data with highly concurrent reads. Should the cache be deleted before updating the database, or the database updated before deleting the cache? These are typical questions.
All four scenarios generally require new source data to reach the destination quickly with minimal delay, while minimizing the impact on the source during synchronization.
Solutions
Scenario One
Logstash’s JDBC plugin can handle the first scenario. It periodically runs SQL queries and sends new data to Elasticsearch. This adds some delay, however, and frequent queries are unnecessary when no new data exists.
With MySQL, Alibaba’s cannal can address those issues. In reality, we may use PostgreSQL and therefore have to choose Logstash.
Scenario Two
For the same database type, such as MySQL, Alibaba’s cannal is still an option. For combinations such as Oracle and PostgreSQL, however, that does not work.
Scenario Three
Database triggers can implement this in any of our existing databases. They add execution time to the original SQL, and adding them changes the relevant table structure. Handling it in application code is another option, but adds complexity.
Scenario Four
Most solutions mix cache-handling logic into business logic, making application developers worry about cache consistency while implementing features. Cache synchronization could be a separate service decoupled from application logic.
None of these existing solutions is ideal.
How Does CDC—Change Data Capture—Solve These Problems?
How CDC Works
Each database has an important log recording all changes:
- Oracle: redo log.
- MySQL: binary log.
- PostgreSQL: write-ahead log.
Software monitors these logs and identifies and extracts data changes when the source database changes. CDC captures detailed insert, update, and delete information, including affected tables, before/after values, and operation timestamps. The changes can be passed to other systems or applications for appropriate processing.
CDC requires no changes to the source database and causes no performance overhead on it.
A Brief Introduction to Kafka Connect and Debezium
Kafka Connect needs two connectors:
- sourceConnector.jar imports source data into a Kafka topic.
- sinkConnector.jar exports topic data to the destination.
Kafka itself does not provide a particularly broad connector selection. This is where Debezium comes in.
Think of Debezium as an implementation of CDC, providing sourceConnector.jar and sinkConnector.jar for many databases.
Scenario Two Example: Oracle-to-PostgreSQL Synchronization Through CDC
For this demo, we need:
- A three-node Kafka cluster.
- Oracle and PostgreSQL installed.
- debezium
- A Docker environment.
All Docker Compose files and configuration instructions are on GitHub. To try it, follow the repository steps in order. Leave a comment under this post if you encounter problems.
Troubleshooting Debezium
The core advice is simple: read the official connector documentation carefully, especially its opening section, How the connector works.
Approaches for Scenarios One and Three
Once you understand CDC and try the source example yourself, handling scenarios one and three is straightforward. They simply use different sourceConnector.jar and sinkConnector.jar packages. The question becomes where to find them.
Search Confluent Hub, which provides many packages, including Elasticsearch source and sink connectors, each with documentation. Alternatively, visit Confluent > Products > Connectors.
For scenario one, Confluent has written a tutorial on synchronizing MySQL data to Elasticsearch, also linked from the Elasticsearch sink connector documentation.
ETL (Extract, Transform, and Load) and Flink CDC
Business scenarios vary and can be complex. We may need more than unchanged copies, particularly for reporting: synchronize different tables, combine their data into a meaningful record during synchronization, then load it into the destination. This is ETL.
Flink CDC is a Flink-based solution. I have not used it myself, but mention it because it may help you make technical decisions. Choose a solution suited to your situation.
Using CDC for Cache Consistency
Our starting assumption is that caching should be a separate service as far as possible, rather than handling it in application code.
First, let us review current caching strategies.
Caching Strategies
Here is ByteByteGo’s diagram:
Cache Consistency on a Single Server
On one server, we generally use faster JVM-level caching. The simplest implementation is a Map; Ehcache is common in practice.
We use these strategies:
- Read-through: the application does not care whether data comes from the cache or database; it simply reads through the cache. The cache decides where to get the data, usually with a loader component. See Ehcache’s read-through definition.
- Write-through: every write goes through the cache. When data is written to it, it also persists the data to the database, with both operations in one transaction, usually using a writer component. See Ehcache’s write-through definition.
Under these strategies, the cache service handles reads and updates. Application developers need not manage them. Cache-aside and write-around, by contrast, require application handling.
How can these strategies be implemented?
Hibernate’s second-level cache implements them.
One of the main challenges of using an application-level cache is ensuring data consistency across entity aggregates. That’s where the second-level cache comes to the rescue
See the Caching section of the Hibernate documentation for configuration.
Cache Consistency in a Distributed System
Suppose Redis stores the cache in a distributed system.
Approach One: Hibernate Second-Level Cache with Redisson
For configuration, see this question and this article.
If using Jedis, search for a suitable solution.
Approach Two: CDC
Monitor database changes with CDC. A cache service simply consumes Kafka data and deletes or updates Redis entries.
The flow becomes:
- Server A updates the database.
- CDC synchronizes the change to Kafka.
- Cache Server consumes the message and deletes or updates Redis.
- Server B now receives the new cached data.
- Both A and B read from Redis, while Cache Server handles updates.
The key practical questions are the CDC service’s stability and latency.
Further reading: Uber’s caching article.
These are solutions for particular caching scenarios. Actual requirements, caching strategy, and team technology choices all need consideration to find a balanced approach. There is no silver bullet.
Our company currently uses:
- Extensive JVM-level caching, with Redis for some scenarios.
- Message queues to synchronize caches across server JVMs.
Closing Thoughts
I have only recently started using CDC for some business scenarios, so I do not yet have strong best-practice experience with Debezium and Kafka Connect to share. This article therefore offers little practical experience to draw on; it mainly introduces CDC concepts, basic examples, and ways to think about different solutions.