一条典型的 AI 后端数据同步链路

数据密集型 AI 后端的标准链路长这样:

1
2
3
4
5
6
7
8
9
10
11
12
13
MySQL(业务主库)

│ binlog

Canal Server(订阅 binlog)

│ Kafka topic

Kafka(消息缓冲)

│ Consumer

Doris / ES(分析 + 搜索)

看起来很标准。但链路越长,能挂的环节越多。生产上我处理过的事故:

  • Canal 订阅位点丢失,从头重做导致下游重复消费 6 小时
  • Kafka 集群重启,Consumer Group rebalance 时部分 partition 没被认领
  • Doris stream load 报 -235 错误(数据格式问题),但 Kafka 消息还在
  • 业务库从主库切到备库,Canal 配置没改导致订阅断了 3 小时

**任何一环出错,最终都是”数据丢了一段时间”或者”数据有重复”**。这篇文章讲清楚这条链路每个环节的故障怎么恢复。

链路一环一环看

第 1 环:MySQL binlog

风险点

  • MySQL 服务重启,binlog 写入位置可能漂移
  • 主从切换后,新主库的 binlog file 名和位点全变了
  • 业务执行了 RESET MASTER,binlog 全部清空

恢复策略

  1. MySQL 必须开启 binlog_row_image=FULL

    1
    SET GLOBAL binlog_row_image = 'FULL';

    缺省 MINIMAL 会丢字段信息,Canal 解析不出完整行。

  2. expire_logs_days 设大一点

    1
    2
    # my.cnf
    expire_logs_days = 7

    给 Canal 重连或重订阅留时间窗口。

  3. 主从切换后 Canal 必须重新配置
    Canal 订阅的是具体主库的 server_id 和 binlog 位点,切主后这俩全变。需要:

    • 修改 Canal instance 配置里的 canal.instance.master.address
    • 重新选择订阅起点(从最新位点,或者从切换时间点)
  4. 永远别 RESET MASTER
    运维误操作的高发点。如果必须清理 binlog,用 PURGE BINARY LOGS TO 'mysql-bin.010' 这种精确清理。

第 2 环:Canal Server

风险点

  • Canal instance 状态丢了(重启后从头开始订阅)
  • Canal 解析 binlog 出错(DDL、特殊字符、JSON 类型)
  • Canal 内存 OOM,订阅位点没持久化

恢复策略

  1. Canal instance 必须配置持久化

    1
    2
    3
    4
    # canal.properties
    canal.instance.tsdb.spring.xml=classpath:spring/tsdb/h2-tsdb.xml
    canal.instance.tsdb.dir=/data/canal/tsdb
    canal.instance.tsdb.url=jdbc:h2:/data/canal/tsdb/canal

    启用 H2 嵌入式数据库持久化位点。重启 Canal 后能从上次位置继续。

  2. 位点丢失后的恢复流程

    • 找到 Canal instance 的位点文件(默认在 conf/{instance}/meta.dat
    • 如果文件坏了,Canal 会从头开始订阅——这意味着下游必须能处理全量回放
    • Doris 端必须用 UNIQUE KEY 模型,ES 端必须用 doc_id 主键,否则会重复
  3. Canal 解析失败的监控

    1
    2
    # Canal 日志
    tail -f logs/canal/canal.log | grep "ERROR"

    常见错误:

    • ERROR com.alibaba.otter.canal.parse.inbound.mysql.MysqlEventParser - dump address /127.0.0.1:3306 has an error, retrying.
    • ERROR - canal parse binlog failed, Table 'xxx' not found in DML
  4. Canal 解析失败的隔离
    配置忽略某些表或错误:

    1
    2
    3
    4
    canal.instance.filter.regex=.*\\..*
    canal.instance.filter.black.regex=^mysql\\.slave_.*$
    canal.instance.tsdb.url=...
    canal.instance.gtidon=false

第 3 环:Kafka 集群

风险点

  • Kafka broker 故障,partition 短暂不可用
  • Consumer Group rebalance 不均衡
  • 消息积压

恢复策略

  1. Kafka 副本数 ≥ 3

    1
    2
    3
    # server.properties
    default.replication.factor=3
    min.insync.replicas=2

    这是数据安全的底线。

  2. Producer 端 acks=all

    1
    props.put(ProducerConfig.ACKS_CONFIG, "all");

    Canal 默认配置已经是 acks=all,但确认下。

  3. Consumer Group rebalance 策略

    1
    2
    # canal kafka sink 配置
    canal.sink.kafka.partitionHash = _key

    按主键 hash 到 partition,保证同一行数据的有序性。

  4. 积压监控

    1
    2
    3
    # 看 Consumer lag
    kafka-consumer-groups.sh --bootstrap-server kafka:9092 \
    --describe --group canal-to-doris

    关键指标:LAG 列。如果 lag > 10000 持续增长,下游消费能力不够。

第 4 环:Doris Stream Load / Routine Load

风险点

  • Stream Load 报错(数据格式不对、字段超长、主键冲突)
  • Routine Load 任务卡住
  • Doris BE 节点故障

恢复策略

  1. Stream Load 错误的兜底

    1
    2
    3
    4
    5
    6
    7
    8
    9
    10
    11
    12
    13
    14
    // 重试 3 次后写入死信 topic
    int retry = 0;
    while (retry < 3) {
    try {
    dorisStreamLoad(record);
    return;
    } catch (Exception e) {
    retry++;
    if (retry >= 3) {
    // 发到死信 topic
    kafkaTemplate.send("rag-dlq", record);
    }
    }
    }
  2. Routine Load 状态监控

    1
    SHOW ALL ROUTINE LOAD\G

    重点关注:

    • State:RUNNING / PAUSED / CANCELLED
    • ErrorLogUrls:失败日志的浏览器 URL
    • DataQualityStats:错误行数 / 总行数
  3. Doris BE 故障
    Doris 是多副本的,单个 BE 挂掉不影响写入。但如果是磁盘满或 OOM,BE 会”卡住”:

    • 检查 BE 节点状态:SHOW BACKENDS\G
    • 查慢查询:SELECT * FROM information_schema.queries ORDER BY duration DESC LIMIT 10;

整条链路的故障恢复剧本

剧本 1:Canal 订阅断了

症状:Kafka 那个 topic 突然没有新消息。

排查

1
2
3
4
5
# 看 Canal 日志
tail -f logs/canal/canal.log

# 看 Canal instance 状态
curl http://canal-admin:8089/api/v1/canal/instance/log?instance=my_instance

恢复

  1. 找到 Canal 报错原因(MySQL 连接失败?位点损坏?)
  2. 如果 MySQL 连接失败:修复 MySQL 网络/权限,重启 Canal instance
  3. 如果位点损坏:把 Canal 位点重置到最后一个健康位置(比如”4 小时前”),让 Canal 重做 4 小时数据
  4. 下游(Doris/ES)必须用 UNIQUE KEY / doc_id 主键去重

剧本 2:Kafka 集群重启

症状:Kafka 短暂不可用后恢复,Consumer lag 持续增长。

排查

1
2
# 看 lag
kafka-consumer-groups.sh --bootstrap-server kafka:9092 --describe --group canal-to-doris

恢复

  1. 检查 Consumer 是否在线(可能 rebalance 出问题)
  2. 重启 Consumer 让它从最新 committed offset 继续
  3. 如果是积压:拉更多 Consumer 并行消费 / 加 Doris 并发

剧本 3:MySQL 主从切换

症状:Canal 日志报”无法连接旧主库”。

恢复

  1. 修改 Canal instance 配置 canal.instance.master.address = new_master
  2. 重启 Canal instance(这会从最新位点开始)
  3. 关键问题:切换前的那部分数据可能没消费完
  4. pt-table-checksum 对账,找缺失的数据窗口期
  5. 手动补数据(或重做 Canal 订阅起点)

剧本 4:Doris BE 故障

症状:Doris 写入报错 -230/-235/-238。

恢复

  1. 确认 BE 状态:SHOW BACKENDS\G
  2. 错误的 -230 通常是 BE 磁盘满 → 清理磁盘空间
  3. -235 是数据格式问题 → 检查上游 JSON 格式
  4. -238 是 tablet 不可用 → 重启 BE 或 add backend
  5. 不要慌着重启整个 Doris,单 BE 故障不影响整体服务

关键的对账机制

任何数据同步链路都必须有对账。光靠 Canal + Kafka + Doris 的位点还不够,因为:

  • Canal 位点可能错位
  • Kafka 消息可能在网络层丢失
  • Doris 写入可能因为 -235 错误被丢

对账脚本(每天跑一次):

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
-- 1. 找出 MySQL 中今天有更新的订单 ID
SELECT order_id FROM orders
WHERE updated_at >= CURDATE();

-- 2. 找出 Doris 中今天写入的订单 ID
SELECT order_id FROM orders_doris
WHERE created_at >= CURDATE();

-- 3. diff
SELECT t1.order_id FROM (
SELECT order_id FROM orders WHERE updated_at >= CURDATE()
) t1
LEFT JOIN (
SELECT order_id FROM orders_doris WHERE created_at >= CURDATE()
) t2 ON t1.order_id = t2.order_id
WHERE t2.order_id IS NULL;

如果 diff 不为空,触发告警 + 自动补数

监控告警清单

指标 阈值 告警级别
Canal instance 状态 ≠ RUNNING 立即 P0
Canal parse 错误 > 0/min P1
Kafka Consumer lag > 10000 P2
Kafka Consumer lag > 100000 P1
Doris stream load 错误率 > 1% P1
Doris BE Alive=false 立即 P0
对账脚本 diff 行数 > 0 P1

我的工具箱

实际项目里我会同时跑几个工具:

  1. maxwelldebezium 作为 Canal 的备选(不同 binlog 解析器,互为冗余)
  2. Kafka Connect 而不是裸 Consumer(自带 offset 管理、错误处理、监控)
  3. Doris 监控用 doris-manager 或自研 exporter(Prometheus + Grafana)
  4. pt-table-checksum 做对账(Percona 出品,老牌工具)
  5. 链路追踪:每条数据从 MySQL 到 Doris 都打一个 trace_id,跨系统串联

意味着什么?

**数据链路的可靠性是”工程问题”,不是”组件问题”**。

  • 单独看 Canal、单独看 Kafka、单独看 Doris,每个都很成熟
  • 但串起来就容易出问题,因为每个组件都有”自己认为对”的语义
  • Canal 认为自己”已经投递到 Kafka”就完事了
  • Kafka 认为自己”已经被 Consumer 拉取”就完事了
  • Doris 认为自己”成功写入”就完事了
  • 真正完整的可靠性是对账 + 监控 + 兜底重放

数据密集型 AI 后端的工程师,最值钱的不是”搭链路”的能力,而是”链断了能快速恢复”的能力。这是运维能力,也是架构能力。

如果你正在搭从 MySQL 到 Doris 的链路,先想清楚:链断了怎么办。比想清楚”链怎么搭”重要 10 倍。

关联阅读