搜索质量取决于索引里的数据是否新鲜、是否完整。把 MySQL、PostgreSQL、MongoDB 与各类 SaaS 的数据同步进 Elasticsearch,是搜索架构的地基。本文覆盖 Connectors 生态、增量同步与 API 写入、CDC 变更捕获与一致性,以及同步冲突处理与运维监控。
1. Connectors 生态
一句话总结: Connectors 是开箱即用的数据源接入框架,把数据库与 SaaS 内容同步成可搜索的 ES 索引。
1.1 Connectors 是什么
Elastic Connectors 是一组原生或服务型连接器,支持 MySQL、PostgreSQL、MongoDB、Oracle 等数据库,以及 SharePoint、Drive、Salesforce 等 SaaS。连接器处理认证、全量+增量抓取、字段映射与幂等写入,把接入成本从「手写同步作业」降为「配置连接器」。
1.2 部署形态
连接器可运行在 Elastic Cloud 托管侧,也可本地部署 Connector Service。本地形态需要一个服务进程与 ES 集群通信,通过连接器 API 管理:
# connector service 配置片段
elasticsearch.host: https://localhost:9200
elasticsearch.api_key: <api-key>
connector_id: mysql-books
1.3 何时选择 Connectors
数据源是常见数据库或 SaaS、同步频率按分钟级即可、不需要强一致事务时,Connectors 是性价比最高的选择;需要毫秒级一致或复杂 ETL 时,走自建同步管道更合适。连接器适合「内容搜索」类接入,不适合「事务一致」类需求。
1.4 一个 MySQL 连接器配置
通过连接器 API 创建数据源与管道,字段映射与过滤规则都在这里声明:
curl -X PUT "localhost:9200/_connector/mysql-books" -H "Content-Type: application/json" -d'
{
"service_type": "mysql",
"index_name": "books",
"pipeline": "books-normalize",
"configuration": {
"host": { "value": "db.internal:3306" },
"database": { "value": "catalog" },
"user": { "value": "reader" },
"tables": { "value": ["books", "authors"] }
}
}
'
配置完成后启动同步作业,连接器按调度周期抓取并写入目标索引,元数据(游标、上次同步时间)自动持久化。
2. 增量同步与 API 写入
一句话总结: 增量同步依赖源端时间戳或游标,API 写入用 bulk 批量提交,两者结合保证同步的高效与可续传。
2.1 增量同步的三类依据
增量同步靠三类依据:更新时间戳(updated_at)、自增主键游标、CDC 日志(binlog/WAL)。时间戳最简单但有同秒更新丢失风险;游标稳定但无法感知删除;CDC 最完整但依赖源端配置。实际生产常把时间戳与游标组合,用主键游标兜底同秒更新:
{
"query": {
"bool": {
"filter": [
{ "range": { "updated_at": { "gt": "2026-10-01T08:00:00Z" } } },
{ "range": { "id": { "gt": 100000 } } }
]
}
}
}
2.2 bulk 批量写入
无论增量还是全量,落 ES 都用 bulk 批量提交,每批几百到几千条,减少请求往返:
curl -X POST "localhost:9200/_bulk" -H "Content-Type: application/json" -d'
{ "index": { "_index": "books", "_id": "1" } }
{ "title": "深入理解搜索引擎", "updated_at": "2026-10-01T08:00:00Z" }
{ "index": { "_index": "books", "_id": "2" } }
{ "title": "Elasticsearch 权威指南", "updated_at": "2026-10-01T08:00:00Z" }
'
2.3 幂等与续传
用源端主键做 ES 文档 _id,重放同一条更新天然幂等。增量任务记录上次同步的游标位置,断点续传从游标继续,避免重复全量。bulk 失败批次要按条解析错误,区分可重试与需人工处理:
# 伪代码:bulk 写入与失败分类
resp = es.bulk(index="books", body=payload, refresh=False)
for item in resp["items"]:
err = item.get("index", {}).get("error")
if err and "retryable" in err.get("type", ""):
retry_queue.put(item) # 网络/限流可重试
elif err:
dead_letter.put(item) # mapping/类型冲突进死信
2.4 增量游标的持久化
同步任务把游标写入状态索引或元数据表,重启后读取续传。状态文件建议单独索引并只保留当前值,避免游标随业务索引的 ILM 被误删。
3. CDC 与一致性
一句话总结: CDC 从源库日志捕获增删改事件,配合版本字段实现近似实时的增量同步与最终一致。
3.1 CDC 的原理
CDC(Change Data Capture)读取数据库 binlog 或 WAL,捕获 insert/update/delete 事件流。相比轮询时间戳,CDC 不放过删除、不丢同秒更新、延迟可到亚秒级。Debezium 是常见的 CDC 中间件,把事件发到 Kafka 再消费写入 ES。
3.2 事件到文档的映射
一条更新事件映射成 ES 的一次 index(upsert),一条删除事件映射成 delete:
{
"op": "u",
"after": { "id": 7, "title": "Kafka 权威指南", "price": 79.9 },
"source": { "ts_ms": 1696000000000 }
}
消费端按 op 决定执行 index/delete,并写入同步版本时间戳供排障。
3.3 最终一致与乱序
CDC 天然是最终一致:事件按提交顺序到达,但重试与分区可能造成短暂乱序。用 _version 或时间戳做乐观并发控制,旧事件不覆盖新事件:
PUT /books/_doc/7
{
"title": "Kafka 权威指南(第 2 版)",
"_version_type": "external",
"_version": 1696000000002
}
3.4 一条 Debezium 链路
Debezium 监听 MySQL binlog,把变更事件写入 Kafka,消费端再落 ES。中间的消息队列提供重放与削峰能力:
{
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"database.hostname": "db.internal",
"database.user": "debezium",
"database.server.name": "catalog-db",
"table.include.list": "catalog.books",
"topic.prefix": "cdc.books"
}
消费端按 Kafka 分区顺序处理事件,把 op=c/u 转 index、op=d 转 delete,并以事件时间戳作为版本号写入,天然解决同分区内的顺序问题。
4. 同步策略
一句话总结: 同步策略决定全量与增量的衔接方式、频率与链路拓扑,是数据新鲜度与成本的平衡点。
4.1 全量 + 增量衔接
首次全量建立基线,随后增量持续跟进;全量期间可能已有增量事件,需要「先增补再切增量」或用版本号合并。全量跑批选低峰,控制并发避免打爆源库。
4.2 同步频率分档
不同业务对新鲜度要求不同:商品库存秒级、价格分钟级、文章标题小时级。同步频率与源库压力、ES 写入量直接相关,按业务价值分档配置,避免全部高频拖垮资源。
4.3 链路拓扑
简单链路:源库 → 同步作业 → ES;高级链路:源库 → Debezium → Kafka → 消费写 ES,中间多一层缓冲,消费失败可重放、可限流。拓扑越简单越好维护,只有需要削峰、多消费者时才引入消息队列。
4.4 频率与批量的工程化
同步频率不宜写死在代码里,放进配置中心按数据源与业务域独立调整。批量大小与源库压力成反比:全量大库用小批量+低并发,增量小变更用大批量+高吞吐。批次按「时间窗口 + 条数上限」双阈值触发,兼顾实时性与吞吐。
| 数据源 | 建议频率 | 批量策略 |
|---|---|---|
| 商品库存 | 秒级 CDC | 小批量高频 |
| 商品详情 | 分钟级 | 中批量低频 |
| 知识库文档 | 小时级 | 全量+增量 |
| 订单流 | 秒级 CDC | Kafka 削峰 |
5. 冲突处理
一句话总结: 冲突来自多源重复、更新竞态与删除漂移,靠统一主键、乐观版本与软删除收敛。
5.1 统一主键
多数据源写入同一索引时,用「源类型 + 源主键」拼成 ES _id,避免不同源的同名 ID 互相覆盖:
curl -X POST "localhost:9200/_bulk" -H "Content-Type: application/json" -d'
{ "index": { "_index": "products", "_id": "mysql_101" } }
{ "title": "无线鼠标", "source": "mysql" }
{ "index": { "_index": "products", "_id": "saas_101" } }
{ "title": "无线鼠标", "source": "saas" }
'
5.2 更新竞态
同一文档被多个同步任务并发更新时,后写覆盖先写。用 external 版本号或 updated_at 比较,只让较新的事件落库;或用版本字段让旧事件静默丢弃。
5.3 删除漂移与软删除
删除事件丢失会导致 ES 残留已删除文档。常见对策:源端软删除(逻辑标记)同步到 ES 过滤;CDC 全链路补齐删除事件;定期对账跑「全量比对删除」清理漂移。对账是最终一致的兜底,周期与数据重要性成正比。
5.4 冲突裁决规则
多源字段冲突要有显式的裁决规则:按数据源优先级、按时间戳最新、或按字段拆分互不覆盖。规则写进同步配置并文档化,遇到争议事件时按规则重放而非人工盲改。
{
"conflict": {
"field_policy": "source_priority",
"priority": ["mysql", "saas", "manual"],
"timestamp_field": "updated_at"
}
}
6. 监控与运维
一句话总结: 同步链路要监控延迟、错误率与积压量,配合重试、告警与对账把故障损失降到最小。
6.1 同步延迟监控
延迟 = 源端事件时间 到 ES 可见时间 的差值。秒级延迟适合监控大盘,超过业务阈值即告警。CDC 链路还要盯 Kafka 消费 lag,积压超过窗口说明消费者吞吐不足。
6.2 错误与重试
同步错误分可重试(网络、限流)与不可重试(字段类型不匹配、mapping 冲突)。可重试按指数退避,不可重试进死信队列并告警。死信保留原始事件便于人工修复后重放。
6.3 对账与数据质量
定期对源库与 ES 做抽样比对,核对文档数、最新时间戳与删除同步。对账发现漂移后触发定向 reindex 或增量补数。同步状态表记录每个连接器的游标、批次号与错误数,作为数据质量的唯一事实来源。
6.4 同步状态表设计
状态表维护每个连接器的运行信息,监控与排障都从它出发:
PUT /sync-status/_doc/mysql-books
{
"last_sync_at": "2026-10-01T08:00:00Z",
"cursor": "id>100000",
"total_indexed": 102400,
"errors_1h": 3,
"dead_letter_count": 2,
"lag_seconds": 5,
"status": "running"
}
对账脚本与告警均查询此表:lag 超阈值告警、errors 持续增长告警、status 非 running 超时告警,形成「状态表驱动」的运维闭环。
7. 实践案例
一句话总结: 电商商品库、知识库文档、订单搜索三类典型同步,覆盖数据库与 SaaS 的主要接入形态。
7.1 电商商品同步
MySQL 商品表 + 库存表用 CDC 同步,标题/类目/价格/库存字段映射进 ES。秒级库存走高频增量,商品详情走分钟级全量+增量,搜索时按库存过滤再排序。
7.2 文档知识库
SharePoint/Drive 用 Connectors 抓取文档,正文解析后进 ES,配合向量字段做语义检索。增量抓取识别文档变更时间戳,删除文档用软删除标记过滤。
7.2 文档知识库
SharePoint/Drive 用 Connectors 抓取文档,正文解析后进 ES,配合向量字段做语义检索。增量抓取识别文档变更时间戳,删除文档用软删除标记过滤。文档同步多走异步解析管道:先抓元数据,正文解析完成后再回填,避免解析阻塞抓取主线程。
7.3 订单搜索
订单量大的场景用 Kafka 削峰,订单事件流消费写入 ES,按状态、金额、时间构建搜索与聚合。索引与业务库解耦,搜索负载不压主库;对账兜底保证搜索库与业务库最终一致:
{
"mappings": {
"properties": {
"order_id": { "type": "keyword" },
"status": { "type": "keyword" },
"amount": { "type": "double" },
"created_at": { "type": "date" }
}
}
}
订单搜索库常用别名承载滚动索引,按日期切分订单索引并用 ILM 治理,查询带时间范围走对应段,既控制规模又保证热数据响应。
8. 总结
| 环节 | 要点 |
|---|---|
| Connectors 生态 | 数据库/SaaS 开箱接入,配置即同步 |
| 增量同步 | 时间戳/游标/CDC,bulk 幂等写入 |
| CDC 一致性 | binlog/WAL 捕获事件,版本号防乱序 |
| 同步策略 | 全量+增量衔接,频率分档、链路拓扑 |
| 冲突处理 | 统一主键、乐观版本、软删除对账 |
| 监控运维 | 延迟/积压/错误监控,死信重试 |
| 实践案例 | 商品、知识库、订单三类落地 |
数据同步是把外部世界搬进搜索索引的管道,核心始终是新鲜度、完整性与一致性的三角平衡。Connectors 降低接入门槛,CDC 提供近似实时的增量,版本号与对账守住一致性,延迟与错误监控让链路可运维。同步链路的设计与容错可阅读《搜索服务架构》,索引与字段规划见《数据建模与 Mapping 设计》,写入性能治理参考《性能调优与缓存策略》。
延伸阅读
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。