MySQL binlog日志机制解析与Canal CDC数据同步实战

MySQL binlog(Binary Log)记录所有数据库表结构和数据变更的操作日志,是 MySQL 主从复制、数据恢复、增量同步的基础设施。Canal 是阿里巴巴开源的 MySQL binlog 增量订阅消费组件,模拟 MySQL Slave 协议实时获取 binlog 事件,广泛应用于缓存更新、搜索引擎同步、数据仓库 ETL 等 CDC(Change Data Capture)场景。

MySQL binlog三种格式与事件类型解析

binlog 有三种格式,各有优缺点:

STATEMENT:记录 SQL 原文。日志量小,但 UUID()、NOW() 等非确定性函数在主从间结果不一致
ROW:记录每行数据的变更前后镜像。数据一致性最好,但日志量大
MIXED:默认使用 STATEMENT,遇到非确定性函数时自动切换为 ROW 格式

查看和修改 binlog 格式:

mysql> SHOW VARIABLES LIKE 'binlog_format';
+---------------+-------+
| Variable_name | Value |
+---------------+-------+
| binlog_format | ROW |
+---------------+-------+
1 row in set (0.00 sec)

mysql> SET GLOBAL binlog_format = 'ROW';
mysql> SET GLOBAL binlog_row_image = 'FULL'; # 记录变更前后完整行数据

binlog 事件类型(ROW 格式下常见的):

WRITE_ROWS_EVENT:INSERT 操作,记录插入的行数据
UPDATE_ROWS_EVENT:UPDATE 操作,记录变更前和变更后的行数据
DELETE_ROWS_EVENT:DELETE 操作,记录被删除的行数据
TABLE_MAP_EVENT:事件前导,映射表名到 table_id
XID_EVENT:事务提交标记,标记一个事务的 binlog 结束位置

使用 mysqlbinlog 工具解析 binlog 文件:

mysqlbinlog --base64-output=DECODE-ROWS -v \
  /var/lib/mysql/mysql-bin.000123 \
  --start-datetime="2026-08-17 00:00:00" \
  --stop-datetime="2026-08-17 12:00:00" \
  --database=order_db

-v 参数将 ROW 格式的二进制行数据解码为伪 SQL(### INSERT INTO `order_db`.`orders`),便于人工审计。

MySQL binlog配置与GTID复制最佳实践

开启 binlog 需要在 my.cnf 中配置:

[mysqld]
server-id = 1
log-bin = /var/lib/mysql/mysql-bin
binlog_format = ROW
binlog_row_image = FULL
binlog_expire_logs_seconds = 604800 # 保留7天
max_binlog_size = 256M

# GTID 配置
gtid_mode = ON
enforce_gtid_consistency = ON
log_slave_updates = ON

# Canal 需要
binlog_checksum = CRC32

GTID(Global Transaction ID)为每个事务分配全局唯一标识,格式为 server_uuid:transaction_id。相比传统的 binlog file + position 方式,GTID 让 Canal 可以基于事务位置精确断点续传,避免主库切换导致位置错乱。

创建 Canal 专用账号并授权:

CREATE USER 'canal'@'%' IDENTIFIED BY 'Canal@2026!';
GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'canal'@'%';
FLUSH PRIVILEGES;

Canal Server部署与实时binlog解析

Canal Server 配置 instance 连接 MySQL:

# conf/example/instance.properties
canal.instance.master.address = 127.0.0.1:3306
canal.instance.dbUsername = canal
canal.instance.dbPassword = Canal@2026!
canal.instance.connectionCharset = UTF-8
canal.instance.filter.regex = order_db\\..*
canal.instance.gtidon = true

使用 Canal Java Client 消费 binlog 事件:

import com.alibaba.otter.canal.client.CanalConnector;
import com.alibaba.otter.canal.client.CanalConnectors;
import com.alibaba.otter.canal.protocol.Message;

CanalConnector connector = CanalConnectors
    .newSingleConnector(new InetSocketAddress("127.0.0.1", 11111),
        "example", "", "");
connector.connect();
connector.subscribe("order_db\\..*");

while (true) {
  Message message = connector.getWithoutAck(1000);
  long batchId = message.getId();
  if (batchId == -1 || message.getEntries().isEmpty()) {
    Thread.sleep(500);
    continue;
  }

  for (Entry entry : message.getEntries()) {
    if (entry.getEntryType() == EntryType.ROWDATA) {
      RowChange rowChange = RowChange.parseFrom(entry.getStoreValue());
      for (RowData rowData : rowChange.getRowDatasList()) {
        EventType eventType = rowChange.getEventType();
        String tableName = entry.getHeader().getTableName();

        switch (eventType) {
          case INSERT:
            handleInsert(tableName, rowData.getAfterColumnsList());
            break;
          case UPDATE:
            handleUpdate(tableName,
                rowData.getBeforeColumnsList(),
                rowData.getAfterColumnsList());
            break;
          case DELETE:
            handleDelete(tableName, rowData.getBeforeColumnsList());
            break;
        }
      }
    }
  }
  connector.ack(batchId);
}

Canal+Kafka数据同步管道与缓存一致性方案

Canal Server 支持直接将 binlog 事件投递到 Kafka:

# canal.properties
canal.serverMode = kafka
canal.mq.servers = kafka:9092
canal.mq.topic = canal-order-db
canal.mq.partition = 0
# 按 table 自动分区,保证同一表的事件顺序
canal.mq.partitionHash = order_db\\..*:id

消费者从 Kafka 读取数据后更新 Redis 缓存,写 Elasticsearch 索引:

from kafka import KafkaConsumer
import json, redis, elasticsearch

consumer = KafkaConsumer(
  'canal-order-db',
  bootstrap_servers=['kafka:9092'],
  group_id='cache-sync-group',
  enable_auto_commit=False
)

r = redis.Redis(host='redis', port=6379)
es = elasticsearch.Elasticsearch(['es:9200'])

for msg in consumer:
  data = json.loads(msg.value)
  table = data['table']
  event_type = data['type']
  
  if event_type == 'INSERT' or event_type == 'UPDATE':
    row = data['data'][0]
    cache_key = f"{table}:{row['id']}"
    r.setex(cache_key, 3600, json.dumps(row))
    es.index(index=table, id=row['id'], body=row)
  elif event_type == 'DELETE':
    row = data['data'][0]
    cache_key = f"{table}:{row['id']}"
    r.delete(cache_key)
    es.delete(index=table, id=row['id'], ignore=[404])
    
  consumer.commit()

通过 partitionHash 配置按主键 ID 哈希分区,同一行数据的 INSERT/UPDATE/DELETE 事件始终进入同一 Kafka partition,保证事件顺序性。消费者以手动方式提交 offset,确保缓存更新成功后才推进消费位点。Canal 记录的 binlog position 持久化在 ZooKeeper/本地文件中,Canal Server 重启后从上次断点继续,不丢不重。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/mysqlbinlog-ri-zhi-ji-zhi-jie-xi-yu-canalcdc-shu-ju-tong-bu/

(0)
小编小编
上一篇 10小时前
下一篇 10小时前

相关推荐