Canal案例

wen java案例 2

本文目录导读:

Canal案例

  1. 核心原理与架构
  2. 实战案例一:缓存双写一致性(Redis 与 MySQL 同步)
  3. 实战案例二:异构数据同步(MySQL 到 Elasticsearch)
  4. 实战案例三:数据变更审计(历史追踪)
  5. 高频踩坑与调优指南
  6. 总结建议

Canal 是阿里巴巴开源的一个基于 MySQL 二进制日志(Binlog)的增量数据订阅和消费组件,它的核心原理是伪装成 MySQL 的从库(Slave),向主库(Master)发送 dump 协议,主库则将 Binlog 推送给 Canal,Canal 解析 Binlog 后,将数据变更以特定格式(如 JSON)推送给下游。

下面从基础原理、核心应用场景、实战案例(含代码)、以及踩坑指南四个维度展开,帮助你深入理解 Canal。


核心原理与架构

  • 主从复制协议:Canal 模拟 MySQL Slave 的交互协议,向 Master 请求 Binlog。
  • Binlog 格式要求:必须设置为 ROW 模式(binlog_format=ROW),这样才能获取到每行数据的变更前后镜像。
  • 组件角色
    • Canal Server:接收 Binlog,解析为事件(Event)。
    • Canal Client:通过 TCP 或 MQ(Kafka/RocketMQ)从 Server 拉取数据,并写入目标存储(如 ES、Redis、HBase)。

实战案例一:缓存双写一致性(Redis 与 MySQL 同步)

业务痛点:高并发下,先更新数据库再删除缓存的策略存在时间窗口,可能导致脏数据,使用 Canal 监听 Binlog 变更,异步更新或删除缓存,可有效避免缓存雪崩和脏读。

场景模拟: 用户表 user 更新了手机号,需要同步更新 Redis 中的用户缓存。

流程图应用更新 MySQL -> Canal 监听 Binlog -> Canal Client 消费变更 -> 应用删除/更新 Redis Key

代码示例(Canal Client 消费核心逻辑)

// 引入依赖(以1.1.5为例,生产建议使用Maven管理)
// com.alibaba.otter:canal.client:1.1.5
public void consumeData() {
    // 1. 连接 Canal Server
    CanalConnector connector = CanalConnectors.newSingleConnector(
        new InetSocketAddress("127.0.0.1", 11111), "example", "", "");
    try {
        connector.connect();
        connector.subscribe("dbName.user"); // 订阅特定库表(正则)
        connector.rollback(); // 回滚到上次消费位点
        while (true) {
            // 2. 批量拉取数据(每次拉取100条,超时500ms)
            Message message = connector.getWithoutAck(100, 500L);
            long batchId = message.getId();
            if (batchId == -1 || message.getEntries().isEmpty()) {
                Thread.sleep(1000);
                continue;
            }
            // 3. 遍历 Binlog 条目
            for (CanalEntry.Entry entry : message.getEntries()) {
                // 只处理数据变更类型
                if (entry.getEntryType() != CanalEntry.EntryType.ROWDATA) continue;
                CanalEntry.RowChange rowChange = CanalEntry.RowChange.parseFrom(entry.getStoreValue());
                CanalEntry.EventType eventType = rowChange.getEventType();
                for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) {
                    // 4. 获取变更后的数据(BeforeColumns / AfterColumns)
                    CanalEntry.Column idColumn = rowData.getAfterColumnsList()
                        .stream()
                        .filter(col -> col.getName().equals("id"))
                        .findFirst()
                        .orElse(null);
                    if (idColumn != null) {
                        // 5. 删除 Redis 缓存(或直接更新)
                        String key = "user:info:" + idColumn.getValue();
                        redisTemplate.delete(key); 
                        System.out.println("缓存已失效: " + key);
                    }
                }
            }
            connector.ack(batchId); // 确认消费成功
        }
    } catch (Exception e) {
        e.printStackTrace();
    } finally {
        connector.disconnect();
    }
}

注意点

  1. 位点管理:消费失败不能调用 ack,否则数据会丢失。
  2. 幂等性:Canal 重启后重新投递,删除 Redis 的操作本身是幂等的,不会造成问题。

实战案例二:异构数据同步(MySQL 到 Elasticsearch)

业务痛点:商品搜索需要复杂的全文检索和多字段组合查询,而 MySQL 的 LIKE '%xxx%' 性能极差。

场景模拟: 商品表 product 价格或库存变化,需要实时同步到 ES 的 product 索引中。

架构升级Canal Server -> Canal Adapter(官方提供的适配器,支持 ES 索引自动映射)或者 Canal Client + 业务逻辑 -> ES RestHighLevelClient

方案对比

  • 方案 A(低代码):使用 Canal Adapter,配置 es.ymlmapping 文件,直接监听表变化,同步到 ES,优点:不用写代码,缺点:复杂逻辑(如字段拼接、多表关联)难以实现。
  • 方案 B(推荐):使用 Canal Client 拉取数据,在业务代码中做 ETL(清洗、转换)后写入 ES。

代码示例(方案 B 核心逻辑)

// 在拉取到 product 表的变动后
if (eventType == CanalEntry.EventType.UPDATE || eventType == CanalEntry.EventType.INSERT) {
    // 1. 构建 ES 文档对象
    Map<String, Object> doc = new HashMap<>();
    for (CanalEntry.Column column : rowData.getAfterColumnsList()) {
        switch (column.getName()) {
            case "product_id":
                doc.put("id", column.getValue());
                break;
            case "product_name":
                doc.put("name", column.getValue());
                break;
            case "price":
                doc.put("price", Double.parseDouble(column.getValue()));
                break;
            // 忽略那些在 ES 中不需要的字段(如 logging_id)
        }
    }
    // 2. 写入 ES
    IndexRequest request = new IndexRequest("product_index")
            .id(String.valueOf(doc.get("id")))
            .source(doc, XContentType.JSON);
    try {
        esClient.index(request, RequestOptions.DEFAULT);
    } catch (IOException e) {
        // 失败时可以投递到 MQ 重试,或者记录日志由定时任务补偿
    }
} else if (eventType == CanalEntry.EventType.DELETE) {
    // 删除 ES 文档
    DeleteRequest request = new DeleteRequest("product_index", 
        rowData.getBeforeColumnsList().get(0).getValue());
    esClient.delete(request, RequestOptions.DEFAULT);
}

实战案例三:数据变更审计(历史追踪)

业务痛点:运营误操作修改了用户余额,需要快速定位“谁在什么时候改了什么”。

实现思路: 将 Binlog 变更记录原样(含变更前、变更后数据)写入 HBase 或 ClickHouse(列式存储,适合分析)。

代码处理核心点

  1. 解析 RowDataBeforeColumnsAfterColumns
  2. 将两个列表的字段名、旧值、新值拼接成 JSON。
  3. 记录操作时间(entry.getHeader().getExecuteTime())和操作类型(eventType)。
// 构造审计日志实体
AuditLog log = new AuditLog();
log.setBizType(entry.getHeader().getSchemaName() + "." + entry.getHeader().getTableName());
log.setEventType(eventType.name());
log.setExecuteTime(new Date(entry.getHeader().getExecuteTime()));
// 记录变更字段差异
for (CanalEntry.Column column : rowData.getAfterColumnsList()) {
    if (column.getUpdated()) { // 判断字段是否被更新
        String fieldName = column.getName();
        String newValue = column.getValue();
        // 查找旧值
        String oldValue = rowData.getBeforeColumnsList().stream()
                .filter(c -> c.getName().equals(fieldName))
                .map(CanalEntry.Column::getValue)
                .findFirst().orElse(null);
        log.addFieldChange(fieldName, oldValue, newValue);
    }
}
auditLogService.save(log);

高频踩坑与调优指南

  1. Binlog 模式必须是 ROW

    • 如果数据库是 STATEMENTMIXED 模式,Canal 拿不到具体的数据行变化,无法正常工作。
    • 排查方法:执行 SHOW VARIABLES LIKE 'binlog_format';,必须为 ROW
  2. MySQL 主从复制用户权限

    • Canal 连接 MySQL 的用户必须有 REPLICATION SLAVEREPLICATION CLIENT 权限。
    • 操作GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'canal'@'%';
  3. 大事务导致 OOM(内存溢出)

    • 如果一个事务更新了 10 万条数据,Canal 会将整个事务的 Binlog 加载到内存解析,如果内存配置不足,Server 会 OOM。
    • 调优:减少 canal.instance.memory.buffer.size 或者调整 JVM 堆内存 -Xms-Xmx
  4. 乱序问题

    • 如果多个线程并发消费并写入目标库,可能会因为网络延迟导致数据乱序(如旧值覆盖新值)。
    • 解决方案:对 主键 HASH 进行分区/分片,保证同一条主键的数据由同一个线程执行(如 Kafka 按表主键分区)。
  5. DDL(数据库结构变更)处理

    • 修改表结构(如新增字段)时,Canal 会发送 DDL 事件,此时你需要判断目标端(如 HBase)的表结构是否也需要更新,否则写入会报错,通常需要建立 DDL 的回调机制进行表结构同步或告警。

总结建议

如果只是简单的 MySQL 与 Elasticsearch/HBase 同步,建议直接使用 Canal Adapter(官方插件)或者 Flink CDC(大数据场景,且有状态管理),如果你是 Java 技术栈,且需要复杂的业务逻辑(如缓存更新策略、特殊字段清洗),直接用 Canal Client 异步消费 会更灵活。

希望这些案例能够帮你将 Canal 落地到实际项目中,如果有具体的报错场景,欢迎带着日志进一步探讨。

抱歉,评论功能暂时关闭!