面试知识库
极高 困难

ES与MySQL数据同步#

一句话答案#

ES 与 MySQL 数据同步推荐 Canal 监听 binlog → Kafka → Consumer 写 ES 的方案,延迟 1-3 秒,保证最终一致性;配合定时全量校验和 alias 切换实现零停机重建索引。

核心要点

同步方案对比#

方案实时性一致性侵入性适用场景
同步双写实时弱(无事务)不推荐
异步双写(MQ)秒级最终一致简单场景
Canal + Kafka秒级最终一致无侵入推荐
定时全量分钟/小时定期一致兜底/冷数据

推荐方案:Canal + Kafka + ES#

MySQL (binlog) → Canal Server → Kafka → ES Consumer → ES

流程详解:
1. MySQL 开启 binlog(ROW 格式)
2. Canal 伪装成 MySQL Slave,拉取 binlog 事件
3. Canal 解析 binlog → 转为结构化消息(INSERT/UPDATE/DELETE)
4. 消息投递到 Kafka(按表名分 Topic)
5. ES Consumer 消费消息:
   - INSERT → ES index 操作
   - UPDATE → ES update 操作(或 index 覆盖)
   - DELETE → ES delete 操作
6. 消费成功后 commit offset
plaintext

Canal 核心原理#

Canal 伪装 MySQL Slave 协议:
1. 向 Master 发送 dump 请求
2. Master 推送 binlog 事件流
3. Canal 解析 ROW 模式 binlog:
   - 包含变更前后的完整行数据
   - 比 STATEMENT 模式更安全(无 SQL 解析歧义)
4. Canal 维护消费位点(binlog filename + position)
5. 支持 failover:HA 模式用 ZooKeeper 选主
plaintext

消费端设计要点#

// 1. 幂等写入(doc_id = MySQL 主键)
IndexRequest request = new IndexRequest("products")
    .id(String.valueOf(row.getId()))  // 相同 id 会覆盖
    .source(buildJson(row));

// 2. 批量写入(提升吞吐)
BulkRequest bulk = new BulkRequest();
for (CanalEntry entry : entries) {
    bulk.add(buildIndexRequest(entry));
}
client.bulk(bulk, RequestOptions.DEFAULT);

// 3. 消费失败重试
// Kafka 消费失败 → 写入死信队列 → 告警 + 人工处理
java

全量同步(兜底 + 初始化)#

场景:新建索引、数据漂移修复

步骤:
1. 创建新索引 products_v2(mapping + settings)
2. 全量 dump MySQL 数据,批量写入 products_v2
3. 增量追赶:全量期间的 binlog 继续消费写入 products_v2
4. 校验数据量一致后,alias 切换:
   POST /_aliases
   { "actions": [
     { "remove": { "index": "products_v1", "alias": "products" }},
     { "add": { "index": "products_v2", "alias": "products" }}
   ]}
5. 删除旧索引 products_v1
plaintext

一致性保障#

问题解决方案
消息丢失Kafka acks=all + 消费 ACK + 死信队列
消息重复ES doc_id = MySQL 主键(天然幂等)
消息乱序按主键 hash 分区保证单行有序
数据漂移定时全量校验(count + 抽样对比)
Canal 宕机HA 模式 + 位点持久化 + 自动恢复

延迟监控#

核心指标:
- binlog → Kafka 延迟(Canal 延迟)
- Kafka → ES 消费延迟(Consumer Lag)
- 端到端延迟 = MySQL 写入时间 - ES 可搜索时间

告警阈值:延迟 > 10s 告警,> 60s 降级(如搜索走 MySQL 兜底)
plaintext
面试回答(2分钟版)

ES 和 MySQL 数据同步推荐 Canal+Kafka 方案。Canal 伪装成 MySQL Slave 拉取 binlog,解析 ROW 格式的变更事件,投递到 Kafka,消费端根据事件类型做 ES 的 index/update/delete 操作。这个方案对业务无侵入,延迟通常 1-3 秒,满足最终一致性。一致性保障靠:Kafka acks=all 防消息丢失、ES doc_id 等于 MySQL 主键保证幂等(重复消费不会出错)、按主键 hash 分区保证单行消息有序。全量同步用于初始化和数据修复:创建新索引 → dump 全量数据 → 增量追赶 → alias 原子切换实现零停机。监控方面重点关注 Canal 延迟和 Consumer Lag,超过阈值告警。同步双写不推荐因为无法保证事务一致性,MySQL 成功但 ES 失败就会数据不一致。

追问与易错

追问方向:

  • “Canal 宕机了怎么办?”→ HA 模式 + ZooKeeper 选主 + binlog position 持久化
  • “怎么处理 MySQL 和 ES 字段不一致?”→ Consumer 做数据转换(宽表拼接/字段映射)
  • “联表数据怎么同步?”→ 监听多表 binlog → 消费时按主表 ID 聚合查询 → 构建 ES 宽文档
  • “零停机重建索引怎么做?”→ 新索引 + 全量写入 + 增量追赶 + alias 切换

易错点:

  • ❌ “同步双写能保证一致”——跨系统无分布式事务,MySQL 成功 ES 失败就不一致
  • ❌ “Canal 能保证消息不丢”——Canal 本身需要 HA 部署 + Kafka 持久化兜底
  • ❌ “binlog STATEMENT 格式就够了”——STATEMENT 有函数不确定性,必须用 ROW 格式