1. 从一次同步延迟事故说起:这不是单纯的代码问题
先交代一下背景。我维护的一个电商搜索服务,SpringBoot 应用,MySQL 作为主库,Elasticsearch 负责商品搜索和类目聚合。业务方时不时会来问"为什么我刚上架的商品搜不到?"、"为什么库存改了前台没反应?"——虽然只是延迟几秒到十几秒的小问题,但赶上大促压测,问题会被无限放大:某次全量同步任务跑了一整夜都没结束,ES 集群 CPU 直接被打满,线上搜索超时率飙到 40%。
那段时间我几乎把和数据同步相关的代码翻了个底朝天,也踩了不少常规文档里根本不会写的坑。这篇就专门聊聊 SpringBoot 环境下,数据库同步 Elasticsearch 这件事的性能优化到底怎么做。标题虽然叫"性能优化",但实际上真正的功夫分布在选型、查询源头、写入端、批次参数、异常补偿这几个环节——任何一个环节掉链子,整体都跑不快。
适合谁来读?如果你正在做或者准备做 MySQL 到 ES 的数据同步,用的是 SpringBoot 技术栈,对 Canal、Logstash 这类重量级组件暂时不想引入,或者已经引入了但发现性能达不到预期,这篇内容应该能给你一些非常具体的参考。后面我会给到可以直接抄的参数配置、代码思路和排查路径,也会解释每个选择背后的理由——这些理由才是你遇到新场景时能做判断的根基。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 同步方案选型:先想清楚你要的是准实时还是最终一致
很多同学一上来就纠结"用什么同步工具",其实第一步应该想清楚的是:你的业务对数据延迟的容忍度到底是多少。这个问题的答案直接决定技术选型,也决定了后续性能优化的天花板。
2.1 三种常见方案的定位差异
市面上主流的 SpringBoot 服务同步方案,本质上逃不出下面这三类:
| 方案 | 实现思路 | 典型延迟 | 复杂度 | 适用场景 |
|---|---|---|---|---|
| 定时全量同步 | 固定时间把整表数据拉到 ES,全量覆盖 | 小时级 | 低 | 数据量小、对实时性无要求、以离线分析为主 |
| 增量轮询同步 | 基于业务表的更新时间戳/自增ID,定时拉取变化数据 | 秒级到分钟级 | 中 | 大多数业务系统的主力方案 |
| Binlog 监听同步 | 通过 Canal 等组件订阅 MySQL binlog,实时解析并写入 ES | 毫秒到秒级 | 高 | 数据变更频繁、业务强依赖准实时搜索 |
我见过不少项目,上来就引入 Canal,理由是"这是最先进的方案"。但实际上团队并没有人能维护 Canal 的 HA 部署和 binlog 位点管理,最后出了问题反而没人能接得住。对于大部分中小规模搜索业务,增量轮询配合定时全量兜底,已经完全够用,而且排查问题非常直观。
2.2 我这边的最终选型
我负责的商品搜索场景,数据量在百万级,变更主要集中在价格、库存、上下架状态这几个字段,业务上"分钟级延迟"完全可接受。所以我的选型是:
- 主链路:增量轮询同步,每 30 秒扫描一次商品表的 update_time 字段,把变更数据同步到 ES。
- 兜底链路:每天凌晨业务低峰期跑一次全量同步,修复增量阶段漏掉的数据。
这个组合的好处是:代码实现简单,出问题容易排查,而且全量任务可以做一些比较激进的性能优化(后面会讲到),不用太担心对线上搜索造成影响。
注意:如果业务真的到了"商品上架必须秒级可搜"的地步,那你确实需要考虑 Canal 或者直接读写穿透(先写 ES 再异步写 MySQL)。但千万不要因为别人用 Canal 你就用,先量化自己的需求再选。
2.3 审视同步性能问题的三个视角
方案定下来之后,性能优化就有了展开的框架。我自己习惯从三个视角去审视同步链路:
- 源头视角:数据库查询是否高效?增量查询能不能避免全表扫描?
- 传输视角:SpringBoot 到 ES 之间的数据包是否过大?批次是否合理?
- 写入视角:ES 端的写入吞吐是否达到瓶颈?索引配置是否为写入场景做了优化?
后面三个章节,就分别对应这三个视角展开。你会发现很多性能问题,根源不在你想的那个环节上——比如 ES 写入慢,可能只是因为数据库查出来的结果集太大,把内存打满了;而数据库查询慢,也可能是因为你没用对增量切分的思路。
3. 源头不卡壳:数据库端查询优化的四个关键点
同步任务本质上是"从数据库读数据,再写到 ES"。先看源头:如果数据库查询这一步慢了,后面再怎么优化写入都是白搭。我在实际问题中总结下来,源头优化有四个最关键的点。
3.1 增量查询必须走索引,否则一次同步就是一次全表扫描
增量同步最基础的写法是:
sql复制SELECT * FROM product WHERE update_time > #{lastMaxUpdateTime} ORDER BY update_time ASC;
很多人以为加了 WHERE 条件就万事大吉,结果 explain 一看,type 是 ALL,全表扫描。update_time 上没有索引,数据量一大必出问题。所以增量查询的第一件事就是给 update_time 字段建索引:
sql复制ALTER TABLE product ADD INDEX idx_update_time (update_time);
建了索引之后,查询性能会有数量级的提升。
但这里还有一个隐藏问题:如果同一秒内有大量数据更新,update_time > 上次记录的最大值 这种写法会把"等于上次最大值"的数据漏掉。这就是典型的边界问题。我自己吃过一次亏:商品批量改价,同一个时间戳里有几千条数据,结果同步完搜索端价格一片混乱。
解决方案有两个:
- 方案 A:把 WHERE 条件改成
update_time >= #{lastMaxUpdateTime},但需要在代码里去重,或者允许重复覆盖(ES 的 doc 覆盖是幂等的,所以重复覆盖没有副作用)。 - 方案 B(我最终采用的):用
update_time结合id做复合游标,即(update_time = #{lastTime} AND id > #{lastId}) OR update_time > #{lastTime},这样既不会漏数据,也不会重复。
3.2 全量同步的切片思路:按主键分片并发读取
如果全量同步直接用 SELECT * FROM product,一次性把几百万行读进 JVM 内存,GC 直接爆炸。而且单线程读+单线程写,全量任务跑几个小时很正常。
正确姿势是按主键做分片,开了多个线程并发读取。具体做法是:
- 查出表的最小主键
MIN(id)和最大主键MAX(id)。 - 根据要开的并发线程数 N,把
[MIN(id), MAX(id)]均分成 N 个区间。 - 每个线程负责读取一个区间内的数据。
比如主键范围是 1~1000000,开 4 个线程,区间就是 [1, 250000]、[250001, 500000]、[500001, 750000]、[750001, 1000000],每个线程只负责自己区间内的 SELECT * FROM product WHERE id BETWEEN ? AND ?。这种方式有两个好处:一是每个线程的查询压力可控,不会造成数据库连接池被瞬间打爆;二是天然支持水平扩展,数据量大就多开几个线程。
3.3 深分页问题必须绕开:用主键游标而不是 OFFSET
很多同步代码里会这么写:
java复制List<Product> list = productMapper.selectPage(pageNum, pageSize);
注意,MyBatis-Plus 的分页最终会翻译成 LIMIT offset, size。当 offset 非常大的时候(比如第 10 万条之后),MySQL 还是要扫描前面所有的行才能定位到目标数据,查询速度断崖式下跌——我在 500 万数据量的表上实测,offset 到 20 万时,一次查询已经需要 3 秒以上。
更糟糕的是,如果数据量大到一定级别,这类查询可能会拖垮整个数据库,影响线上业务。
正确做法是主键游标方式:
sql复制SELECT * FROM product WHERE id > #{lastId} ORDER BY id ASC LIMIT 1000;
每次取完一批,记录这批最大的 id,下一批用这个 id 作为起点继续取。这个思路和上面增量的复合游标是一脉相承的,它把查询成本从"扫描 offset 行再返回"变成了"直接走主键索引跳转到目标位置",性能稳定且不随页数增长而劣化。
3.4 数据库连接池和查询超时的平衡
同步任务通常会并发读取数据库,如果配置不当,同步任务本身就可能把连接池占满,导致线上业务查询阻塞。我自己踩过一次:连接池最大连接数设了 20,而全量同步开了 12 个线程,每个线程还有多表关联查询,瞬间占掉了大部分连接,结果线上商品详情接口出现了缓慢。
我现在固定使用 HikariCP 做连接池,配置上比之前严谨很多:
yaml复制spring:
datasource:
hikari:
maximum-pool-size: 30
minimum-idle: 10
connection-timeout: 3000
validation-timeout: 1500
max-lifetime: 1800000
其中 maximum-pool-size 不是越大越好,需要根据数据库的 max_connections 来倒推。比如数据库允许 200 个连接,线上服务占用 120 个,同步服务最多就分配 50 个,留出至少 30 个余量,防止服务重启或其他突发情况。
另外,同步查询要单独设置查询超时。我的做法是在 Mapper 接口上用 @Options(timeout = 30) 控制单条 SQL 的执行时间,超过 30 秒直接报错,由上层重试机制接管。为什么需要这个?因为同步任务跑在深夜,一旦某条 SQL 因为锁等待被卡住,重试可能瞬间堆积成雪崩。超时是一种快速失败机制,给重试争取时间。
4. 写入端吞不下:ES 写入瓶颈的定位与解决
数据库读出来只是第一步,能不能快速写到 ES 才是真正的考验。我当时全量同步跑了一夜没完成,根子就在写入端。
4.1 单条写入是最容易犯的错
很多第一版同步代码长这样:
java复制for (Product product : productList) {
IndexRequest request = new IndexRequest("product")
.id(String.valueOf(product.getId()))
.source(objectMapper.writeValueAsString(product), XContentType.JSON);
restHighLevelClient.index(request, RequestOptions.DEFAULT);
}
一条一条地发 HTTP 请求。这在数据量小(几千条以内)时没感觉,但一旦数据量上来,问题立刻暴露:
- 每条数据一次网络往返,开销极高。
- ES 每收到一条写入请求,都可能触发一次 refresh,产生大量小 segment。
- 并发高时,ES 端线程池被打满,响应变慢,客户端堆积重试,雪崩。
我实测过 10 万条数据,单条写入耗时大约 8~10 分钟,而用 bulk 批量写入只需要 15 秒左右——差了 40 倍不止。所以单条写入这个写法,在新老项目中都应该直接杜绝。
4.2 Bulk 批量写入的正确使用姿势
ES 官方提供的 Bulk API 就是用来解决批量写入的。核心是用一个请求发送多条数据,减少网络往返。但批量大小怎么定,是有讲究的。
业内推荐的经验值是:批次大小在 1000~5000 条,或者数据总大小在 5MB~15MB 之间,具体看哪个条件先达到。
为什么是这个范围?Bulk 请求太小,没法充分利用批处理优势;Bulk 请求太大,ES 端一次性解析和索引大量文档,JVM 堆内存压力陡增,反而容易触发 OOM 或 GC 长时间停顿。我自己压测的结果是:单批次 5000 条(每条约 1KB,总大小约 5MB)时,吞吐量达到峰值,再往上加批次大小,吞吐量反而下降。
随着 ES 版本升级,现在官方推荐的 bulk 大小标准倾向于拿字节数做判断。我目前的实现方案是:以条数为主限制,同时用字节数做兜底:
java复制public List<BulkOperation> buildBulkOperations(List<Product> list) {
List<BulkOperation> operations = new ArrayList<>();
int batchBytes = 0;
for (Product p : list) {
byte[] bytes = objectMapper.writeValueAsBytes(p);
batchBytes += bytes.length;
operations.add(IndexOperation.of(op -> op
.index("product")
.id(String.valueOf(p.getId()))
.document(p)));
if (operations.size() >= 5000 || batchBytes >= 8 * 1024 * 1024) {
// 达到阈值,触发一次批量提交
executeBulk(operations);
operations.clear();
batchBytes = 0;
}
}
// 提交剩余数据
if (!operations.isEmpty()) {
executeBulk(operations);
}
return operations;
}
注意一个细节:批量操作里尽量让每条文档的字节数规格化,避免某一条特别大的文档拖垮整个批次。如果业务里有大文本字段(比如商品详情),可以考虑在构建索引时单独处理——要么把大字段从 ES 索引里去掉,要么把大字段改成 store: false,让 ES 只索引不存储,这能显著降低写入压力。
4.3 写入线程池:并发不是越大越好
Bulk 提交是 I/O 密集操作,可以考虑多线程并发提交。但线程数需要根据 ES 集群的分片数和硬件配置来确定。
比较实用的估算公式是:最优并发线程数 ≈ ES 集群数据节点数 × 每个节点的分片数 × 0.5~0.75。
比如 3 个数据节点,每个节点 5 个分片,最优并发度大约在 7~11 之间。我之前曾一股脑把线程池开到 50,结果 ES 集群 CPU 和磁盘 I/O 全部飙满,bulk 请求大量返回 429 和 503,吞吐量反而不如 10 个线程的时候。后来我把线程池定在 10,并用 Semaphore 做了信号量限流,保证任何时候最多只有 10 个 bulk 请求在途:
java复制ExecutorService executor = new ThreadPoolExecutor(
10, 10,
0L, TimeUnit.MILLISECONDS,
new LinkedBlockingQueue<>(1000),
new ThreadFactoryBuilder().setNameFormat("es-bulk-sync-%d").build(),
new ThreadPoolExecutor.CallerRunsPolicy()
);
CallerRunsPolicy 这个拒绝策略很关键:当阻塞队列满时,它不会丢弃任务,而是让调用线程自己执行,相当于天然限流反压。我自己在这个场景里没有用 DiscardPolicy(会丢数据),也没有用 AbortPolicy(会抛异常导致同步直接失败)。
4.4 索引导数优化:导入期间关掉你不需要的东西
ES 写入过程中有两个参数对写入吞吐影响极大:refresh_interval 和副本数。
先说说 refresh 机制。ES 默认每秒 refresh 一次,每次 refresh 会把缓冲区的数据生成一个新的 segment——频繁的 refresh 会产生大量小 segment,后续 merge 会消耗大量 CPU。对搜索实时性要求不那么高的场景,批量导入期间可以把 refresh_interval 调大到 30s:
bash复制PUT /product/_settings
{
"index": {
"refresh_interval": "30s"
}
}
全量导入完成后,再把它调回默认的 1s。这个操作在 500 万数据导入场景下,实测可以把写入吞吐提升约 30%~40%。
再来说副本。ES 写入时默认需要把数据复制到副本分片才算完成,这相当于写入了双倍的数据量。批量导入期间,可以先把副本数临时调整为 0:
bash复制PUT /product/_settings
{
"index": {
"number_of_replicas": 0
}
}
导入完成后再恢复副本数,ES 会自动把副本数据补上。这个"先关副本、再开副本"的思路,在数据量大的全量同步中属于标配操作。
注意事项:副本为 0 期间,如果所在节点宕机,这部分数据是会丢的。所以这个操作只适合在可重新导入、可容忍短时数据缺失的场景(比如凌晨全量任务)使用。增量在线同步期间,一定要保持正常副本数。
5. 一次全量同步卡死事故的完整排查链路:从 CPU 100% 到 14 倍性能提升
方法论讲了不少,但真正值钱的往往是完整走一遍问题的排查链路。这一节我用真实事故案例,完整还原我是如何一步步定位并修复"全量同步跑了十几个小时还完不成"这个问题的。
5.1 事故现场与初步判断
某天早上发现,昨晚 23 点启动的全量同步任务还在运行。看监控面板,ES 集群 CPU 使用率已接近 100%,但每秒写入条数却在持续走低,从最初的上万条掉到了几百条。与此同时,MySQL 的 CPU 使用率也很高,商品表所在实例的慢查询日志暴涨。
我第一反应是 ES 写入慢,但看了 ES 监控,发现写入线程池队列其实并没有堆积太多,说明 ES 并没有真正被打满——那问题大概率在"读"这一侧。
5.2 排查链路第一站:MySQL 慢查询日志
慢查询日志里出现频次最高的是一条深分页查询:
sql复制SELECT * FROM product
WHERE status = 1
ORDER BY id ASC
LIMIT 3500000, 1000;
offset 已经到 350 万行,这条查询每执行一次要扫描 350 万行主键索引,耗时 4 秒以上。这就是我前面提到的深分页问题,线上全量同步代码最开始用了 MyBatis-Plus 的简单分页,数据量一大就劣化到不可用。
这里也说明了一个排查规律:同步任务性能问题不要先怀疑 ES,先看数据是从哪儿来的。我把数据源头的慢查询解决掉,往往写入侧的瓶颈也就解开了。
5.3 排查链路第二站:发现批处理和线程池配置的问题
我继续看同步服务的日志,发现一个问题:Bulk 请求平均耗时非常高,而且频繁出现超时重试。但奇怪的是,单条 bulk 的文档数并不多,只有几百条——说明批次大小没配好时,ES 处理大量小 bulk 请求也会产生额外开销。
再往下查,发现同步服务配置了 50 个线程并发提交 bulk。50 个线程同时打向 3 节点 ES,每个节点瞬间收到大量小请求,导致 ES 的单个分片同时处理的请求数过多,写入性能被反压。这正是我前面说的"并发不是越大越好"。
5.4 修复方案与效果对比
定位到三个问题后,我做了一次集中修复:
- 读取侧:把深分页改成主键游标(id > lastId)方式,配合 4 线程分片并发读取。
- 批次侧:把 bulk 大小调整为单批次 5000 条,同时用总字节数兜底(8MB)。
- 并发侧:线程池从 50 压缩到 10,并加上信号量限流。
- 索引侧:全量导入期间把 refresh_interval 调整到 30s、副本数临时置为 0。
修复前:全量 500 万数据,跑了 14 个小时还没跑完,任务最终被杀掉重新跑。
修复后:同样的数据量,耗时约 52 分钟,近 14 倍左右的性能提升,ES 集群 CPU 稳定在 50% 左右,MySQL 慢查询清零。
那之后我把这套"主键游标 + 分片并发 + 批量写入 + 索引调参"的组合固化成了一套统一的同步模板,增量同步和全量同步共用同一套底层逻辑,差别只在数据切分方式和任务触发方式。
6. 稳定运行期的边界条件与一致性兜底:防止性能优化带来的坑
性能调优只是第一步,更考验人的是在长期运行中保证数据一致性和稳定。这里分享几个我实际踩过、也最终解决的边界问题。
6.1 删除数据同步不上:软删除是必须的
数据库里的记录被物理删除后,增量同步基于 update_time 是搜不到这条数据的,ES 里那条文档就成了"孤魂野鬼"——前端搜索还能搜出已经删掉的商品。最直接的解法是改为软删除:业务表加一个 deleted 字段(或 status 字段),删除时更新为已删除,同步任务捕捉到该变更后,在 ES 端执行 delete 操作。
如果业务量实在无法改造为软删除,那就必须在全量同步阶段做一次整体比对,把 ES 中已存在但 DB 中已不存在的文档删除。这个比对成本随数据量线性增长,所以能软删除就尽量软删除,从机制上避免这个问题。
6.2 时间精度与更新覆盖问题
MySQL 的 DATETIME 默认精度是秒。如果同一秒内同一条记录被更新了两次,增量同步的游标记录的是"秒级时间戳+最后一条 ID",第二次更新就可能被漏掉。我遇到的实际场景是:运营后台批量改价格,同一时刻把 A 商品从 100 改成 99,又从 99 改成 98,结果 ES 这边只同步到了 99。
这块的兜底策略我目前用两层:
- 把游标推进和业务读取分开。读取时先按上个周期的"最大更新时间和对应 ID"过滤,但处理成功后,游标位置向后推进时以下一批最大更新时间为准,而不是以本批数据的处理完成为准。这样即便本批数据在写入时有失败,重试时也能覆盖到边界。
- 配合定时对账:每隔 15 分钟跑一次"最近 1 小时变更数据重同步"任务,做幂等覆盖。由于 ES 的写入按主键 ID 是天然幂等的,重跑对线上搜索不会有负面影响。
6.3 大批量导入期间对线上搜索的影响
全量同步虽然放在凌晨,但哪怕是低峰期,ES 集群的资源也是有限的。如果你在导入期间完全不限制写入速度,有可能把 ES 集群的搜索延迟打上去——第二天早上业务方反馈"搜索变慢了",就很尴尬。
我的做法是给全量任务加一个"动态限速"逻辑:每秒检测 ES 集群的搜索延迟和写入队列长度,当指标超过阈值时,自动降低并发写入线程的令牌桶速率。等到指标恢复正常,再逐渐提升速率。这个思路并不复杂,本质上就是基于监控数据的自适应反馈控制,但它能保证同步任务永远不抢业务流量。
6.4 重试与失败补偿:队列版还是定时版
同步任务在运行过程中,可能会遇到 ES 集群临时不可用、网络抖动、单条数据 JSON 序列化失败等问题。如果失败直接抛出异常,整个任务回滚,代价太大;如果静默吞掉,又会造成数据不一致。
我的设计是引入了两级重试:
- 第一级:bulk 提交失败时,把失败的请求单独收集,延迟 2 秒、5 秒、10 秒做指数退避重试,最多 3 次。
- 第二级:3 次重试仍然失败的文档,写入一个"失败补偿表"(MySQL 中一张专门记录同步失败任务的表),由定时任务每 5 分钟扫描补偿表,重新提交到 ES。每条失败记录最多重试 10 次,超过 10 次就发送告警通知人工介入。
这套机制极大降低了同步任务的"彻底失败"概率。实际运行中,绝大多数故障都是瞬时抖动,第一级重试就能解决;真到了第二级,说明 ES 集群确实出了不小的问题,及时告警反而能帮团队提前发现隐患。
7. 我现在的固定动作:一套可以直接落地的优化清单
踩过这些坑之后,现在我接手任何"SpringBoot 同步数据到 Elasticsearch"项目,都会按固定的一套清单做优化,这里直接分享给你。
源头侧:
update_time必须有索引;复合游标(时间+ID)替代update_time > lastTime。- 全量同步严禁
LIMIT OFFSET分页,一律采用主键游标或按主键区间分片。 - 同步服务连接池单独配置,最大连接数 = 数据库 max_connections × 30%。
写入侧:
- 批量写入使用 Bulk API,单批次 5000 条或 8MB 字节数兜底。
- 写入线程池为核心线程数 10,阻塞队列 1000,拒绝策略用 CallerRunsPolicy。
- 大批量导入前调整索引:
refresh_interval=30s、副本数临时为 0,完成后恢复。 - 动态限速:写入速率基于 ES 集群搜索延迟自适应调整。
一致性侧:
- 业务表一律软删除,同步任务捕捉到软删除状态后在 ES 端执行 delete。
- 先推进业务读取再推进游标,配合定时对账做幂等覆盖。
- 两级重试机制:bulk 失败先指数退避重试,重试失败进补偿表,定时任务扫描补偿,超限告警。
监控侧(这个容易被忽略,但非常重要):
- 同步任务要埋点记录:每次任务读取行数、bulk 成功条数、bulk 失败条数、任务耗时、ES 集群写入队列长度。
- 这些指标打进 Prometheus + Grafana,设置阈值告警。如果同步任务连续 N 分钟写 0 条,立刻告警。
我自己建了一个简单的同步指标记录,用 Micrometer 接入 Prometheus:
java复制@Bean
public MeterRegistryCustomizer<MeterRegistry> metricsConfig() {
return registry -> registry.config().commonTags("application", "db-es-sync");
}
// 在每次 bulk 成功后记录
Counter.builder("es.sync.bulk.success")
.register(meterRegistry)
.increment(bulkSize);
这看起来只是运维层面的事,但性能优化最怕的就是没有数据——你觉得自己优化了,但说不清到底快了多少,也说不清哪个环节是瓶颈。有了埋点之后,每次优化都能用数字说话。
这次先分享到这里。说句实在话,数据库同步 Elasticsearch 这件事本身不复杂,容易出问题的从来都是那些不起眼的细节:索引有没有建、游标怎么设计、批次会不会太大、副本要不要临时关、失败了怎么兜底。把这些细节都理清楚,你的同步任务不仅跑得快,而且还稳。希望上面的排查链路和参数配置能给你一些参考。
