Skip to content

搜索与存储同步 ​

#数据同步 · #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 的区别 ​

维度CanalDebezium
支持数据库MySQLMySQL/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")

相关专题 ​


维度CanalDebeziumFlink CDC
架构独立服务Kafka Connect 插件Flink 作业
数据源MySQLMySQL/PG/Mongo/OracleMySQL/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 集群)
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 --> DW
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 改为 1sP99 -4s
增加并行度Kafka 分区数 = 消费者数吞吐 ×N
减少序列化Protobuf 替代 JSONCPU -30%
本地缓存关联数据避免消费时查 DBP99 -50ms
异步写入ES async bulk吞吐 ×3
批注模式

💬 文章评论

暂无评论,来说点什么吧 👇

编程学习笔记