基于 Debezium + Kafka 实现 CDC

2026/9/4 18:30:00Henry2 阅读0 点赞2 评论

为了使得下游能够及时获取好友关系等上游状态,不必每次同步 RPC 查询,可以采用 CDC 对数据库变更进行捕获并下发至 Kafka 等中间件,供下游服务消费。

其核心思想:

  • 数据库是关系状态的事实源
  • Debezium 直接消费 MySQL binlog
  • Kafka 承载 CDC 事件
  • 下游服务消费 Kafka 事件并更新本地 / Redis 分布式缓存
  • 下游获取此类状态优先读取本地 / Redis 分布式缓存,避免频繁请求某个服务

相比业务层主动发送 Kafka Event,CDC 的优势:​业务代码无需额外保证写 DB + 发事件的一致性​。

本文采用 Docker 简易化部署 Debezium + Kafka。

MySQL 前置条件

Debezium MySQL CDC 依赖 binlog:

CONF
[mysqld]
server-id=1
log_bin=mysql-bin
binlog_format=ROW
binlog_row_image=FULL

CDC 用户至少需要下列权限,这边已经提前创建好了用户 debeziumuser

SQL
GRANT SELECT, RELOAD, SHOW DATABASES,
      REPLICATION SLAVE, REPLICATION CLIENT
ON *.*
TO 'debeziumuser'@'%';

Kafka 部署及配置

下方的配置可以简化,因为我还有 K8s 的业务,本文中并未涉及,可忽略。

Kafka 监听只需关注容器内网络的配置即可:PLAINTEXT://:19092 -> PLAINTEXT://kafka:19092

SH
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 核心配置

SH
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:

JSON
{
  "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 的连接信息,此外 namedatabase.server.id 要注意每个 connector 需要不同。

关键配置:

  • topic.prefix:决定了 CDC 下发到 Kafka 的 Topic 前缀;
  • include.schema.changes:是否捕获数据库 schema 的变化;
  • database.include.list:捕获哪些数据库的变更;
  • table.include.list:捕获对应数据库中哪些表的变更(需要加上数据库前缀)。

添加完成后记得查询一下状态:

JSON
{
    "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 时,消息大致如下:

JSON
{
  "schema": {},
  "payload": {
    "before": {},
    "after": {},
    "op": "u"
  }
}

常见 op

TEXT
r = snapshot read
c = create
u = update
d = delete

根据这个基本结构,我们可以在 Golang 工程 Infra 层中进行简单封装:

GO
/*
 * @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 反序列化,这样不同业务域可以用同一套事件结构。

GO
/*
 * @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 格式:

TEXT
<topic.prefix>.<database>.<table>

直接消费该 Topic 中的消息即可,对于 Consumer GroupID,这里有几种情况:

  • 第一种是每个下游实例是本地维护缓存,这种情况下,每个实例的 Consumer GroupID 应该保证不同,分别消费一次;
  • 第二种则是使用的 Redis 等分布式缓存,这种情况下则将各实例消费者设置为相同的 GroupID,避免重复消费。

YuelaiGroup 字节星球

评论区

  • 陶小桃Blog
    #1
    陶小桃Blog2026/9/4 20:56:50

    看上去牛逼,但是我不懂,路过

    • Henry
      #1
      Henry2026/9/4 22:30:11
      个人认证YuelaiGroup, Software Engineer