搜索与存储同步
#数据同步 · #Canal · #Debezium · #Elasticsearch · #MySQL · #CDC
业务系统中,MySQL 负责事务写入,Elasticsearch 负责复杂查询。如何保证两者数据一致是常见挑战。本专题系统梳理数据同步的方案、架构与工程实践。
为什么需要数据同步
- MySQL 擅长:事务、关联查询、精确匹配
- ES 擅长:全文检索、聚合分析、模糊匹配
- 业务需求:写入走 MySQL 保证事务性,查询走 ES 保证性能
两者数据需要保持一致,但不可能做到强一致,目标是最终一致性。
同步方案对比
| 方案 | 实时性 | 一致性 | 侵入性 | 复杂度 |
|---|---|---|---|---|
| 同步双写 | 实时 | 弱 | 高 | 低 |
| 异步双写 | 秒级 | 中 | 高 | 中 |
| Binlog 监听 | 秒级 | 强 | 无 | 中 |
| 定时全量同步 | 分钟级 | 强 | 无 | 低 |
同步双写
方案
python
def create_product(product):
# 写 MySQL
db.insert(product)
# 写 ES
es.index(index="products", id=product.id, body=product.to_dict())问题
- MySQL 成功但 ES 失败 → 数据不一致
- 增加主链路延迟
- ES 不可用时影响写入
改进:异步双写
python
def create_product(product):
db.insert(product)
# 异步发送到消息队列
mq.send("product_sync", {"action": "create", "data": product.to_dict()})
# 消费者异步写入 ES
def sync_to_es(msg):
action = msg["action"]
data = msg["data"]
if action == "create":
es.index(index="products", id=data["id"], body=data)
elif action == "update":
es.update(index="products", id=data["id"], body={"doc": data})
elif action == "delete":
es.delete(index="products", id=data["id"])基于 Binlog 的 CDC 同步
架构
MySQL (Binlog) → Canal/Debezium → Kafka → 同步服务 → Elasticsearch
→ Redis
→ 其他下游优势
- 零侵入:不修改业务代码
- 实时性:毫秒级延迟
- 可靠性:基于 Binlog 位点,支持断点续传
- 通用性:一份 Binlog 可供多个下游消费
Canal
原理
Canal 伪装为 MySQL Slave,接收 Master 的 Binlog 事件。
MySQL Master → Binlog → Canal Server → Canal Client → 下游系统
(伪装 Slave)部署配置
properties
# canal.properties
canal.id = 1
canal.ip =
canal.port = 11111
canal.destinations = example
# instance.properties
canal.instance.master.address = 127.0.0.1:3306
canal.instance.dbUsername = canal
canal.instance.dbPassword = canal
canal.instance.filter.regex = mydb\\..*消费示例
java
// Canal Client 消费 Binlog
CanalConnector connector = CanalConnectors.newSingleConnector(
new InetSocketAddress("127.0.0.1", 11111), "example", "", "");
connector.connect();
connector.subscribe("mydb\\.products");
while (true) {
Message message = connector.getWithoutAck(100);
long batchId = message.getId();
if (batchId != -1) {
for (Entry entry : message.getEntries()) {
RowChange rowChange = RowChange.parseFrom(entry.getStoreValue());
for (RowData rowData : rowChange.getRowDatasList()) {
switch (rowChange.getEventType()) {
case INSERT:
syncToES("create", rowData.getAfterColumnsList());
break;
case UPDATE:
syncToES("update", rowData.getAfterColumnsList());
break;
case DELETE:
syncToES("delete", rowData.getBeforeColumnsList());
break;
}
}
}
connector.ack(batchId);
}
}Debezium
与 Canal 的区别
| 维度 | Canal | Debezium |
|---|---|---|
| 支持数据库 | MySQL | MySQL/PostgreSQL/MongoDB/Oracle/... |
| 部署方式 | 独立服务 | Kafka Connect 插件 |
| 输出 | 自定义协议 | 标准 Kafka Topic |
| 生态 | 阿里系 | 开源社区,Kafka 生态 |
| Schema 管理 | 无 | 支持 Schema Registry |
Kafka Connect 配置
json
{
"name": "mysql-connector",
"config": {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"database.hostname": "mysql",
"database.port": "3306",
"database.user": "debezium",
"database.password": "dbz",
"database.server.id": "184054",
"topic.prefix": "dbserver1",
"database.include.list": "mydb",
"table.include.list": "mydb.products,mydb.orders",
"schema.history.internal.kafka.bootstrap.servers": "kafka:9092",
"schema.history.internal.kafka.topic": "schema-changes.mydb"
}
}消息格式
json
{
"before": {"id": 1, "name": "旧名称", "price": 99},
"after": {"id": 1, "name": "新名称", "price": 199},
"source": {
"version": "2.4.0",
"connector": "mysql",
"ts_ms": 1689012345000,
"db": "mydb",
"table": "products"
},
"op": "u",
"ts_ms": 1689012345123
}同步服务设计
数据转换
MySQL 表结构和 ES 索引结构往往不同,需要转换层:
python
# MySQL 行数据 → ES 文档
def transform_product(row_data):
"""将 MySQL 行转换为 ES 文档"""
doc = {
"id": row_data["id"],
"name": row_data["name"],
"description": row_data["description"],
"price": float(row_data["price"]),
"category_path": row_data["category_name"], # 可能需要关联查询
"tags": row_data["tags"].split(",") if row_data["tags"] else [],
"created_at": row_data["created_at"].isoformat(),
"search_text": f"{row_data['name']} {row_data['description']}", # 组合搜索字段
}
return doc批量写入
python
from elasticsearch.helpers import bulk
# 批量同步,提升性能
def batch_sync_to_es(events, batch_size=500):
actions = []
for event in events:
if event["op"] in ("c", "u"): # create/update
doc = transform_product(event["after"])
actions.append({
"_index": "products",
"_id": doc["id"],
"_source": doc,
})
elif event["op"] == "d": # delete
actions.append({
"_index": "products",
"_id": event["before"]["id"],
"_op_type": "delete",
})
if len(actions) >= batch_size:
bulk(es_client, actions)
actions = []
if actions:
bulk(es_client, actions)一致性保障
延迟补偿
python
# 定时对账任务:比对 MySQL 和 ES 数据
def reconciliation_check():
# 查询最近 1 小时修改的数据
mysql_data = db.query("""
SELECT id, updated_at FROM products
WHERE updated_at > NOW() - INTERVAL 1 HOUR
""")
for row in mysql_data:
es_doc = es.get(index="products", id=row["id"], ignore=404)
if not es_doc or es_doc["_source"]["updated_at"] != row["updated_at"]:
# 数据不一致,重新同步
full_data = db.query("SELECT * FROM products WHERE id = ?", row["id"])
es.index(index="products", id=row["id"], body=transform_product(full_data))
log.warn(f"数据不一致已修复: product_id={row['id']}")全量重建
当数据严重不一致或索引结构变更时,需要全量重建:
python
def full_rebuild(index_name):
# 1. 创建新索引(带版本号)
new_index = f"{index_name}_v2"
es.indices.create(index=new_index, body=mapping)
# 2. 全量导入
offset = 0
batch_size = 1000
while True:
rows = db.query(f"SELECT * FROM products LIMIT {batch_size} OFFSET {offset}")
if not rows:
break
actions = [{"_index": new_index, "_id": r["id"], "_source": transform_product(r)} for r in rows]
bulk(es_client, actions)
offset += batch_size
# 3. 切换别名(零停机)
es.indices.update_aliases(body={
"actions": [
{"remove": {"index": f"{index_name}_v1", "alias": index_name}},
{"add": {"index": new_index, "alias": index_name}},
]
})
# 4. 删除旧索引
es.indices.delete(index=f"{index_name}_v1")常见问题与解决
顺序问题
同一条记录的多次变更必须按顺序同步:
- Kafka 中按主键分区,保证同一记录的事件有序
- 消费端单线程处理同一分区
关联数据同步
宽表场景(ES 文档包含多表数据):
python
# 商品表变更 → 直接同步
# 分类表变更 → 需要同步所有关联商品
def on_category_change(category_id, new_name):
# 查询该分类下所有商品
products = db.query("SELECT id FROM products WHERE category_id = ?", category_id)
# 批量更新 ES 中的分类名称
for product in products:
es.update(index="products", id=product["id"],
body={"doc": {"category_name": new_name}})同步延迟监控
python
# 监控同步延迟
def monitor_sync_lag():
# 获取 Binlog 最新位点
latest_binlog_pos = get_mysql_binlog_position()
# 获取消费者当前位点
consumer_pos = get_consumer_position()
lag = latest_binlog_pos - consumer_pos
if lag > THRESHOLD:
alert(f"数据同步延迟过大: {lag} events")相关专题
- Elasticsearch — 搜索引擎原理
- MySQL — Binlog 与复制
- Kafka — 消息传输层
- 消息驱动架构 — 异步架构模式
- 缓存技术 — 缓存一致性方案
Flink CDC — 新一代数据同步方案
Canal vs Debezium vs Flink CDC
| 维度 | Canal | Debezium | Flink CDC |
|---|---|---|---|
| 架构 | 独立服务 | Kafka Connect 插件 | Flink 作业 |
| 数据源 | MySQL | MySQL/PG/Mongo/Oracle | MySQL/PG/Mongo/Oracle |
| 输出 | 自定义协议 | Kafka Topic | 任意 Sink(ES/Redis/Kafka/DB) |
| 全量+增量 | 需手动全量 | ✅ Snapshot + Streaming | ✅ 无锁 Snapshot + Streaming |
| 数据转换 | 需额外服务 | 需 Kafka Streams/Flink | ✅ 内置 SQL/DataStream |
| Exactly Once | ❌ | ✅ (Kafka 事务) | ✅ (Checkpoint) |
| 多表 Join | ❌ | ❌ (需额外处理) | ✅ (Flink SQL Join) |
| Schema 变更 | 需重启 | 自动感知 | 自动感知 |
| 运维复杂度 | 低 | 中 | 高(需 Flink 集群) |
Flink CDC 架构
mermaid
flowchart TB
subgraph Source["数据源"]
MySQL["MySQL<br/>(Binlog)"]
PG["PostgreSQL<br/>(WAL)"]
end
subgraph Flink["Flink 集群"]
CDC["CDC Source<br/>(无锁快照 + 增量)"]
Transform["数据转换<br/>(Flink SQL / DataStream)"]
Sink["Sink"]
end
subgraph Target["目标"]
ES["Elasticsearch"]
Redis["Redis"]
Kafka["Kafka"]
DW["数据仓库"]
end
MySQL -->|"Binlog"| CDC
PG -->|"WAL"| CDC
CDC --> Transform
Transform --> Sink
Sink --> ES
Sink --> Redis
Sink --> Kafka
Sink --> DWFlink CDC SQL 示例
sql
-- 1. 定义 MySQL CDC Source
CREATE TABLE products_source (
id INT,
name STRING,
price DECIMAL(10, 2),
category_id INT,
updated_at TIMESTAMP(3),
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'hostname' = 'mysql-host',
'port' = '3306',
'username' = 'flink',
'password' = 'xxx',
'database-name' = 'mydb',
'table-name' = 'products'
);
-- 2. 定义 ES Sink
CREATE TABLE products_es (
id INT,
name STRING,
price DECIMAL(10, 2),
category_name STRING,
updated_at TIMESTAMP(3),
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'elasticsearch-7',
'hosts' = 'http://es-host:9200',
'index' = 'products'
);
-- 3. 多表 Join + 写入 ES(Canal/Debezium 做不到的!)
INSERT INTO products_es
SELECT p.id, p.name, p.price, c.name AS category_name, p.updated_at
FROM products_source p
LEFT JOIN categories_source c ON p.category_id = c.id;Flink CDC 的核心优势:一个作业同时完成"数据采集 + 数据转换 + 数据写入",无需 Canal→Kafka→消费者→ES 的长链路。多表 Join 在 Flink 内部完成,避免了消费端关联查询的复杂性。
同步延迟 SLA 设计
延迟分层
端到端同步延迟 = Binlog 延迟 + 传输延迟 + 处理延迟 + 写入延迟
┌──────────────────────────────────────────────────────────┐
│ 阶段 │ 典型延迟 │ 瓶颈点 │
├───────────────┼─────────────┼────────────────────────────┤
│ Binlog 产生 │ < 1ms │ MySQL 主从延迟 │
│ CDC 采集 │ 10-100ms │ 网络 RTT + 解析 │
│ 消息传输 │ 5-50ms │ Kafka/Flink 内部 │
│ 数据转换 │ 1-10ms │ 计算复杂度 │
│ 写入目标 │ 10-100ms │ ES bulk / Redis pipeline │
├───────────────┼─────────────┼────────────────────────────┤
│ 总计 (P50) │ 50-200ms │ │
│ 总计 (P99) │ 500ms-2s │ 批量写入 + 网络抖动 │
└──────────────────────────────────────────────────────────┘SLA 定义
| 级别 | P99 延迟目标 | 适用场景 | 方案 |
|---|---|---|---|
| L1 (实时) | < 1s | 搜索结果、商品价格 | Flink CDC 直写 |
| L2 (准实时) | < 10s | 用户画像、推荐 | Canal/Debezium + Kafka |
| L3 (近实时) | < 1min | 报表、统计 | 批量同步 + 定时对账 |
| L4 (离线) | < 1h | 数据仓库、归档 | 定时全量 + 增量 |
延迟监控方案
python
# 方案 1: 心跳表法(最准确)
# 在 MySQL 中创建心跳表,定时写入时间戳
# 消费端收到心跳消息后,计算延迟
# MySQL 端:每秒写入心跳
# INSERT INTO sync_heartbeat (id, ts) VALUES (1, NOW(3))
# ON DUPLICATE KEY UPDATE ts = NOW(3);
# 消费端:计算延迟
def check_sync_lag():
# 从 ES 中读取心跳记录
es_heartbeat = es.get(index="sync_heartbeat", id=1)
es_ts = parse_timestamp(es_heartbeat["ts"])
# 当前时间 - ES 中的心跳时间 = 端到端延迟
lag = time.now() - es_ts
metrics.gauge("sync_lag_seconds", lag.total_seconds())
if lag > SLA_THRESHOLD:
alert(f"同步延迟超过 SLA: {lag}")promql
# Prometheus 监控同步延迟
# 方案 2: 基于 Kafka Consumer Lag 估算
kafka_consumergroup_lag{group="sync-service"} * avg_message_interval_seconds
# 告警规则
- alert: SyncLagExceedsSLA
expr: sync_lag_seconds > 10
for: 2m
labels:
severity: critical
annotations:
summary: "数据同步延迟超过 SLA ({{ $value }}s)"延迟优化手段
| 优化点 | 方法 | 效果 |
|---|---|---|
| 减少批量等待 | ES bulk 从 5s 改为 1s | P99 -4s |
| 增加并行度 | Kafka 分区数 = 消费者数 | 吞吐 ×N |
| 减少序列化 | Protobuf 替代 JSON | CPU -30% |
| 本地缓存关联数据 | 避免消费时查 DB | P99 -50ms |
| 异步写入 | ES async bulk | 吞吐 ×3 |
登录后即可发表评论 👇