基于 Debezium + Kafka 实现 CDC
为了使得下游能够及时获取好友关系等上游状态,不必每次同步 RPC 查询,可以采用 CDC 对数据库变更进行捕获并下发至 Kafka 等中间件,供下游服务消费。
其核心思想:
- 数据库是关系状态的事实源
- Debezium 直接消费 MySQL binlog
- Kafka 承载 CDC 事件
- 下游服务消费 Kafka 事件并更新本地 / Redis 分布式缓存
- 下游获取此类状态优先读取本地 / Redis 分布式缓存,避免频繁请求某个服务
相比业务层主动发送 Kafka Event,CDC 的优势:业务代码无需额外保证写 DB + 发事件的一致性。
本文采用 Docker 简易化部署 Debezium + Kafka。
MySQL 前置条件
Debezium MySQL CDC 依赖 binlog:
[mysqld]
server-id=1
log_bin=mysql-bin
binlog_format=ROW
binlog_row_image=FULL
CDC 用户至少需要下列权限,这边已经提前创建好了用户 debeziumuser:
GRANT SELECT, RELOAD, SHOW DATABASES,
REPLICATION SLAVE, REPLICATION CLIENT
ON *.*
TO 'debeziumuser'@'%';
Kafka 部署及配置
下方的配置可以简化,因为我还有 K8s 的业务,本文中并未涉及,可忽略。
Kafka 监听只需关注容器内网络的配置即可:PLAINTEXT://:19092 -> PLAINTEXT://kafka:19092
docker run -d \
--name kafka \
--network docker-net \
--hostname kafka \
--restart unless-stopped \
-p 9092:9092 \
-p 29092:29092 \
-v kafka-data:/var/lib/kafka/data \
-v kafka-secrets:/etc/kafka/secrets \
-v kafka-config:/mnt/shared/config \
-e KAFKA_NODE_ID=0 \
-e KAFKA_PROCESS_ROLES=broker,controller \
-e KAFKA_LISTENERS='CONTROLLER://:9093,PLAINTEXT://:19092,PLAINTEXT_HOST://:9092,K8S://:29092' \
-e KAFKA_ADVERTISED_LISTENERS='PLAINTEXT://kafka:19092,PLAINTEXT_HOST://127.0.0.1:9092,K8S://host.docker.internal:29092' \
-e KAFKA_LISTENER_SECURITY_PROTOCOL_MAP='CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT,K8S:PLAINTEXT' \
-e KAFKA_INTER_BROKER_LISTENER_NAME=PLAINTEXT \
-e KAFKA_CONTROLLER_QUORUM_VOTERS='0@kafka:9093' \
-e KAFKA_CONTROLLER_LISTENER_NAMES=CONTROLLER \
-e KAFKA_LOG_DIRS=/var/lib/kafka/data \
-e KAFKA_LOG_RETENTION_MS=3600000 \
-e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1 \
-e KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR=1 \
-e KAFKA_TRANSACTION_STATE_LOG_MIN_ISR=1 \
-e KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS=0 \
-e KAFKA_HEAP_OPTS='-Xmx512m -Xms256m' \
apache/kafka:4.3.0
Debezium 核心配置
docker run -d \
--name debezium-connect \
--network docker-net \
--hostname debezium-connect \
--restart unless-stopped \
-p 8083:8083 \
-e BOOTSTRAP_SERVERS='kafka:19092' \
-e GROUP_ID='debezium-connect' \
-e CONFIG_STORAGE_TOPIC='debezium_connect_configs' \
-e OFFSET_STORAGE_TOPIC='debezium_connect_offsets' \
-e STATUS_STORAGE_TOPIC='debezium_connect_statuses' \
-e CONNECT_CONFIG_STORAGE_REPLICATION_FACTOR=1 \
-e CONNECT_OFFSET_STORAGE_REPLICATION_FACTOR=1 \
-e CONNECT_STATUS_STORAGE_REPLICATION_FACTOR=1 \
-e HEAP_OPTS='-Xmx512m -Xms256m' \
quay.io/debezium/connect:3.6
启动后需要对 Debezium Connect 进行配置,这里先罗列一下常用配置接口:
| 操作 | 接口 |
|---|---|
| 查看所有 Connector | GET /connectors |
| 创建 Connector | POST /connectors |
| 查看状态 | GET /connectors/{name}/status |
| 查看配置 | GET /connectors/{name}/config |
| 更新配置 | PUT /connectors/{name}/config |
| 重启 Connector | POST /connectors/{name}/restart |
| 删除 Connector | DELETE /connectors/{name} |
首先创建一个 connector:
{
"name": "interaction-cdc",
"config": {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"database.hostname": "mysql8",
"database.port": "3306",
"database.user": "debeziumuser",
"database.password": "debeziumpassword",
"database.server.id": "23333",
"topic.prefix": "cdc_test",
"database.include.list": "yuelaiengine",
"schema.history.internal.kafka.bootstrap.servers": "kafka:19092",
"schema.history.internal.kafka.topic": "schemahistory.fullfillment",
"include.schema.changes": "false",
"table.include.list":"yuelaiengine.relationship"
}
}
这里面需要关注数据库和 Kafka 的连接信息,此外 name 和 database.server.id 要注意每个 connector 需要不同。
关键配置:
topic.prefix:决定了 CDC 下发到 Kafka 的 Topic 前缀;include.schema.changes:是否捕获数据库 schema 的变化;database.include.list:捕获哪些数据库的变更;table.include.list:捕获对应数据库中哪些表的变更(需要加上数据库前缀)。
添加完成后记得查询一下状态:
{
"name": "interaction-cdc",
"connector": {
"state": "RUNNING",
"worker_id": "192.168.107.3:8083",
"version": "3.6.2.Final"
},
"tasks": [
{
"id": 0,
"state": "RUNNING",
"worker_id": "192.168.107.3:8083",
"version": "3.6.2.Final"
}
],
"type": "source"
}
消费 CDC 事件
上述配置完成之后,一旦对应数据在 MySQL 中发生变更,就会及时被 Debezium 捕获并下发 Kafa 事件。后续就应该解析事件结构并进行消费。
Kafka 事件结构
Kafka Connect JSON Converter 开启 schema 时,消息大致如下:
{
"schema": {},
"payload": {
"before": {},
"after": {},
"op": "u"
}
}
常见 op:
r = snapshot read
c = create
u = update
d = delete
根据这个基本结构,我们可以在 Golang 工程 Infra 层中进行简单封装:
/*
* @Author: Henry csthenry@foxmail.com
* @Date: 2026-09-03 15:12:41
* @LastEditors: Henry csthenry@foxmail.com
* @LastEditTime: 2026-09-03 15:12:43
* @FilePath: /engine/internal/infra/cdc/debezium/model.go
* @Description:
*
* Copyright (c) 2026 by Henry email: csthenry@foxmail.com, All Rights Reserved.
*/
package debezium
type Operation string
const (
OperationRead Operation = "r"
OperationCreate Operation = "c"
OperationUpdate Operation = "u"
OperationDelete Operation = "d"
)
type Message[T any] struct {
Payload *Payload[T] `json:"payload"`
}
type Payload[T any] struct {
Before *T `json:"before"`
After *T `json:"after"`
Operation Operation `json:"op"`
}
这里我忽略掉了基本用不上的 schema 字段,下面是支持泛型的 JSON 反序列化,这样不同业务域可以用同一套事件结构。
/*
* @Author: Henry csthenry@foxmail.com
* @Date: 2026-09-03 16:01:41
* @LastEditors: Henry csthenry@foxmail.com
* @LastEditTime: 2026-09-03 16:16:09
* @FilePath: /engine/internal/infra/cdc/debezium/decode.go
* @Description:
*
* Copyright (c) 2026 by Henry email: csthenry@foxmail.com, All Rights Reserved.
*/
package debezium
import (
"encoding/json"
"errors"
"fmt"
)
var ErrTombstone = errors.New("debezium tombstone")
// Decode 解析 debezium kafka 消息 payload
func Decode[T any](data []byte) (*Payload[T], error) {
if len(data) == 0 {
return nil, ErrTombstone
}
var message Message[T]
if err := json.Unmarshal(data, &message); err != nil {
return nil, fmt.Errorf("decode debezium message failed: %w", err)
}
return message.Payload, nil
}
https://developer.aliyun.com/ask/584158
在 MySQL 的 Binlog 中,删除操作不会直接包含被删除的数据,而是包含一个表示删除操作的事件。当 CDC 工具捕获到这个删除事件时,它通常会生成一个tombstone消息,这是一个特殊的空消息或者包含极少信息的消息,用于指示某个键值对已被删除。
这里注意要特殊处理一下 Debezium Tombstone,这类情况需要单独处理或忽略。
消费 Kafka
具体消费的 Code 在这里就不罗列了,无非就是调用 Kafka 的 SDK,然后按照上述代码对具体业务的事件进行解析,并管理好 Consumer Goroutine 的生命周期即可。
这里说明一下该 CDC 工具生成的 Topic 格式:
<topic.prefix>.<database>.<table>
直接消费该 Topic 中的消息即可,对于 Consumer GroupID,这里有几种情况:
- 第一种是每个下游实例是本地维护缓存,这种情况下,每个实例的 Consumer GroupID 应该保证不同,分别消费一次;
- 第二种则是使用的 Redis 等分布式缓存,这种情况下则将各实例消费者设置为相同的 GroupID,避免重复消费。
YuelaiGroup 字节星球
评论区
看上去牛逼,但是我不懂,路过
北京WindowsChrome
@陶小桃Blog 我造!过气网站居然有人光顾!欢迎👏
福建省macOSChrome