Debezium CDC 变更数据捕获入门与实战
Mei Lin | 2026-09-02T01:08:00 | DevOps, Database
介绍 Debezium 的 CDC 原理和架构,演示如何通过 MySQL Binlog 捕获数据变更事件,实现数据库同步、缓存更新和事件驱动架构。
# Debezium CDC 变更数据捕获入门与实战 ## 什么是 CDC CDC(Change Data Capture)是一种捕获数据库变更事件的技术。与定时轮询不同,CDC 通过读取数据库的事务日志(如 MySQL Binlog)来实时捕获每一行数据的增删改操作。 Debezium 是最流行的开源 CDC 平台,基于 Kafka Connect 构建。 ## 架构 ``` MySQL (Binlog) | Debezium Connector (Kafka Connect) | Kafka Topics (每张表一个 topic) | +-- Consumer A: 同步到 Elasticsearch +-- Consumer B: 更新 Redis 缓存 +-- Consumer C: 触发业务事件 ``` ## 环境搭建 ```yaml # docker-compose.yml version: "3.8" services: mysql: image: mysql:8.0 environment: MYSQL_ROOT_PASSWORD: root command: --server-id=1 --log-bin=mysql-bin --binlog-format=ROW ports: - "3306:3306" zookeeper: image: confluentinc/cp-zookeeper:7.5.0 environment: ZOOKEEPER_CLIENT_PORT: 2181 kafka: image: confluentinc/cp-kafka:7.5.0 depends_on: - zookeeper environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 ports: - "9092:9092" connect: image: debezium/connect:2.7 depends_on: - kafka - mysql environment: BOOTSTRAP_SERVERS: kafka:9092 GROUP_ID: 1 CONFIG_STORAGE_TOPIC: connect_configs OFFSET_STORAGE_TOPIC: connect_offsets STATUS_STORAGE_TOPIC: connect_statuses ports: - "8083:8083" ``` ## 注册 Connector ```bash curl -X POST http://localhost:8083/connectors \ -H "Content-Type: application/json" \ -d '{ "name": "mysql-connector", "config": { "connector.class": "io.debezium.connector.mysql.MySqlConnector", "tasks.max": "1", "database.hostname": "mysql", "database.port": "3306", "database.user": "debezium", "database.password": "dbz_password", "database.server.id": "184054", "topic.prefix": "dbserver1", "database.include.list": "mydb", "table.include.list": "mydb.orders,mydb.users", "schema.history.internal.kafka.bootstrap.servers": "kafka:9092", "schema.history.internal.kafka.topic": "schema-changes" } }' ``` ## 事件格式 每个变更事件的结构: ```json { "before": { "id": 1001, "status": "pending", "amount": 99.99 }, "after": { "id": 1001, "status": "completed", "amount": 99.99 }, "source": { "version": "2.7.0", "connector": "mysql", "name": "dbserver1", "ts_ms": 1700000000000, "db": "mydb", "table": "orders", "server_id": 1, "file": "mysql-bin.000003", "pos": 154 }, "op": "u", "ts_ms": 1700000000123 } ``` - `op`: 操作类型 — `c`(create), `u`(update), `d`(delete), `r`(read/snapshot) - `before`: 变更前的数据(UPDATE/DELETE 有值) - `after`: 变更后的数据(CREATE/UPDATE 有值) ## 消费变更事件 ```java @KafkaListener(topics = "dbserver1.mydb.orders") public void handleOrderChange(ConsumerRecord record) { JsonNode event = objectMapper.readTree(record.value()); String op = event.get("payload").get("op").asText(); JsonNode after = event.get("payload").get("after"); switch (op) { case "c": log.info("New order: {}", after.get("id")); searchService.index(after); break; case "u": log.info("Order updated: {}", after.get("id")); cacheService.invalidate("order:" + after.get("id")); searchService.update(after); break; case "d": JsonNode before = event.get("payload").get("before"); log.info("Order deleted: {}", before.get("id")); cacheService.delete("order:" + before.get("id")); searchService.delete(before.get("id").asLong()); break; } } ``` ## 常见应用场景 1. **数据库同步**:MySQL -> Elasticsearch/MongoDB 2. **缓存失效**:数据变更时自动更新 Redis 缓存 3. **审计日志**:记录所有数据变更历史 4. **事件驱动微服务**:替代应用层发布事件,保证一致性 5. **数据湖入湖**:实时将变更流入数据湖 ## 生产注意事项 1. **MySQL 配置**:必须开启 ROW 格式的 Binlog 2. **Binlog 保留期**:至少保留 3 天,给 Connector 故障恢复留余量 3. **Schema 变更**:Debezium 自动处理,但 `ALTER TABLE` 可能需要重新快照 4. **大表初始快照**:首次启动会做全表快照,大表需注意内存和时间 5. **监控指标**:关注 Connector 的 lag 和状态 ## 总结 Debezium CDC 让你用事件驱动的方式响应数据变更,比应用层双写更可靠(无一致性问题),比定时轮询更实时。结合 Kafka 生态,它是构建实时数据管道和事件驱动架构的利器。