
做个人项目的时候最烦的不是写CRUD而是要把业务数据放到ES里去做搜索和分析。最初我都是手工写脚本一条SQL查出来然后循环写入ES后来发现脚本越来越多代码重复、配置乱、跑起来还要担心中途失败。某个周末我干脆花了半天时间做了个个人数据同步ES的小工具把所有同步任务统一管起来。这篇文章把设计思路和实际踩坑记录下来给遇到类似问题的朋友一个参考。这个小工具解决的典型问题很简单业务数据在MySQL或者MongoDB里但搜索、聚合要从Elasticsearch里查两边数据需要保持一致。对于个人博客、订单管理、商品demo这类项目数据量不会太大但一旦手动同步很快就会遇到三个痛点全量重建太麻烦、增量更新容易漏、任务失败没有日志。小工具的目标就是把这些都收拢成一条命令。如果你也想自己做一个轻量级的同步工具这篇文章会把需求拆解、方案选型、核心实现、问题排查都讲清楚有些代码可以直接拿去改。1. 工具定位与需求拆解1.1 为什么个人项目需要把数据同步到ES我先说下我的使用场景。个人项目里我习惯用MySQL存业务数据比如文章表、评论表、商品表因为关系型数据库做事务和关联查询都很顺手。但一旦要支持全文检索、多字段过滤、聚合统计MySQL的like查询和普通索引就很吃力了特别是几个字段联合搜索的时候SQL越写越复杂性能还差得离谱。这时候Elasticsearch的价值就出来了。ES基于倒排索引天生适合文本搜索和聚合分析还能在查询时做“分词、高亮、加权、范围过滤”。我只要在MySQL里维护好数据再定期把数据同步到ES搜索接口完全走ES体验会好很多。但这带来一个新问题数据同步怎么做才能不掉链子。如果你只是同步一次那写个一次性脚本也就算了。关键是数据会持续变化今天新增一篇文章明天修改一个价格后天删除一个分类如果同步机制不可靠ES里的数据和数据库就会越差越多。等到你想起来去核对时只能重新全量同步一次下次还是会乱。这个小工具的出发点就是把“全量同步”和“增量同步”都固化下来让它自己能循环跑。1.2 小工具需要解决的核心问题动手之前我列了几个必须满足的要求。个人工具最怕的是“为了做一个工具而做一个工具”功能看着很多真正用起来全是坑。第一配置化。不能每次新增一张表就写一遍代码。我希望在properties或者YAML里加一段配置指定数据源、同步SQL、目标索引然后工具就能自动跑起来。这样后面加新的同步任务只需要改配置不需要重新编译打包。第二全量同步和增量同步都要支持。全量同步用于首次建立索引或者数据结构调整后重建增量同步用于定时轮询把新增和修改的数据及时反映到ES里。第三任务可观察。同步工具跑起来是个后台任务如果没有日志、没有状态记录出了问题你会完全不知道它停在哪里。所以我要求每个任务都有明确的日志输出并且同步断点要持久化进程重启后能从上次进度继续。第四轻量。我个人项目没有条件搞一套大数据组件也不希望引入消息队列、中间件这些重量级依赖。一个小工具能跑在单机或一台小服务器上打成jar包就能运行这是底线。基于这些要求我大概画出了一个工具边界数据源连接、查询读取、ES批量写入、任务调度、同步状态存储。后面所有开发都是围绕这几块展开的。2. 方案选型为什么最终自己动手写2.1 先看看开源同步工具能不能直接搬在决定自己写之前我把主流的开源同步方案都过了一遍包括DataX、MongoShake、Canal、Debezium还有SQL Server的CDC/CT方案。每个都有各自的优势但对个人项目来说都显得有些“重”。工具/方案适用场景优点缺点DataX离线批量同步支持MySQL、ES、HDFS等成熟稳定字段映射丰富增量同步需要自己拼where条件没有断点续传管理MongoShakeMongoDB实例间同步也可以同步到ES专为MongoDB设计支持增量部署和配置较重个人小数据量有点杀鸡用牛刀Canal/Debezium基于数据库binlog/WAL的CDC同步实时准确能捕获删除操作依赖消息队列或额外组件配置复杂度高SQL Server CDC/CTSQL Server的变更捕获机制数据库原生支持增量仅限SQL Server需要开启数据库功能我的数据源主要是MySQL也可能会加MongoDB。如果用DataX单次全量同步确实很快但程序化和定时化不方便。MongoShake主要面向MongoDB如果我还想同步MySQL就得再找一套工具。Canal和Debezium实时性很好但要处理binlog解析、消息中间件、消费位点保存对个人项目来说运维成本太高了。最后我给自己定的方案是核心同步逻辑自己写使用JDBC分页读取数据使用ES官网客户端批量写入。这样既能满足全量和增量又能把配置控制在自己手里灵活度最高。2.2 同步策略如何取舍数据同步通常有三种策略我分别想了它们的适用边界。全量同步最简单就是把数据源当前所有符合条件的数据全部同步到ES。优点是实现简单适合初次建索引和数据修复缺点是数据量大时耗时长每次都全量跑一遍显然不现实。增量同步按业务时间字段来做比如每次查询update_time 上次同步时间的数据。实现也不难还能定时反复跑。它的弱点在于对“删除”无能为力因为物理删除的记录不会出现在查询结果里。我的处理方式是尽量使用软删除标记如果确实需要物理删除同步再考虑接入CDC方案。还有一种更接近实时的CDC方案比如MySQL开启binlog通过Canal或者Debezium订阅变更事件。这种方式最精准能够捕获insert、update、delete但个人环境里开启binlog、维护消费位点都挺折腾。我决定把时间和精力先放在前两种策略上把接口预留好将来需要再接入CDC。在具体使用上我采用“首次全量 定时增量 出错后手动全量恢复”的组合。平时定时任务跑增量数据出问题时重新跑一次全量两个策略都能在小工具里一键触发这就满足了我的日常需求。2.3 技术栈和核心库选择开发语言我选了Java。虽然Python写脚本更快但Java在连接池、线程管理、类型转换、打包部署上更适合做成一个长期运行的工具。具体依赖如下Spring Boot负责任务调度和配置管理HikariCP作为数据库连接池JDBC Template或MyBatis处理查询官方elasticsearch-java客户端负责ES写入。这里特别说一下为什么不用RestHighLevelClient。旧版HighLevelClient确实用得多但从ES 8开始官方推荐新的Java ClientAPI更清晰性能也不错。新客户端底层使用JSON和HTTP和BulkProcessor搭配得很好。如果项目基于Spring Boot 3建议直接用co.elastic.clients:elasticsearch-java。我用Spring的Scheduled来做定时任务没有上Quartz。个人工具的任务数量有限Scheduled配合cron表达式已经完全够用。如果以后任务多了需要分布式调度再替换成XXL-Job也不迟。3. 核心实现同步任务抽象与实现细节3.1 任务的抽象模型小工具最核心的设计是抽象出“数据源”和“同步任务”两个概念。数据源是一组连接配置比如MySQL的地址、用户名、密码同步任务描述的是“从哪个数据源查什么数据写到哪个ES索引”。一个SyncTask大致包含这些字段数据源引用、查询SQL、ES索引名、业务主键字段、增量字段、同步批次大小、ES刷新间隔配置。data-sources: mysql-blog: url: jdbc:mysql://localhost:3306/blog username: root password: ****** sync-tasks: - name: sync_blog_articles source: mysql-blog querySql: SELECT id, title, content, category_id, update_time FROM article WHERE update_time :lastSyncTime esIndex: articles idColumn: id incrColumn: update_time batchSize: 1000通过这个抽象每次新增一张表只需要在YAML里增加一个task配置代码完全不用动。不管是同步文章表还是商品表底层逻辑都是“读取配置 - 分页查询 - 批量写入ES - 更新断点状态”。3.2 全量同步的落地要点全量同步最忌讳的是把数据一次性全捞进内存。我早期写脚本时用一条不带LIMIT的SQL查出来几十万条数据直接内存溢出。后来改成“分页查询批量写入”之后问题就消失了。分页查询有两种方式。传统limit offset分页简单但offset越大数据库扫描越慢个人项目数据量小时可用。我更推荐基于排序字段的keyset分页比如按id排序每次只查比当前游标大的前N条这样即使数据再多性能也稳定。下面这段代码是典型的keyset分页逻辑long lastId 0L; while (true) { ListArticle rows jdbcTemplate.query( SELECT id, title, content, category_id, update_time FROM article WHERE id ? ORDER BY id ASC LIMIT ?, new Object[]{lastId, batchSize}, (rs, rowNum) - new Article( rs.getLong(id), rs.getString(title), rs.getString(content), rs.getLong(category_id), rs.getTimestamp(update_time) ) ); if (rows.isEmpty()) { break; } // 将rows转换为Bulk请求提交到BulkProcessor indexRows(rows); lastId rows.get(rows.size() - 1).getId(); }为什么用id而不是其他字段做游标因为id通常稳定且唯一不会因为数据更新而变化。如果同步的全量查询有复杂的过滤条件也可以使用where条件里的业务字段做游标但必须保证游标字段有索引并且值不重复。3.3 增量同步与断点续传增量同步的核心是回答一个问题上次同步到哪儿了我把同步状态保存在一个sync_task_state表里每次增量任务开始前读取last_sync_time任务执行结束后用新的时间更新这张表。这样进程重启也不会丢进度。CREATE TABLE sync_task_state ( task_name VARCHAR(64) PRIMARY KEY, last_sync_time DATETIME, last_sync_batch INT, update_time DATETIME );实现增量查询时要注意边界条件。比如任务在10:00:00执行查询update_time 2025-01-01 10:00:00但在10:00:00.5秒有一条数据被更新了update_time更新为10:00:00正好等于记录下的上次时间下一轮同步就会漏掉它。我的处理方式是做一个小偏移每次记录上次同步时间时往前挪几秒比如last_sync_time - 2秒这样即使边界处有数据变动下一轮也会再查一遍。配合ES写入时使用业务主键作为_id重复写入不会产生重复数据这个策略非常实用。3.4 批量异步写入ES的性能调优ES写入性能的关键在Bulk。逐条写入document会非常慢因为每次都要建立连接、处理请求、刷新分片。BulkProcessor可以自动攒批达到一定条数、大小或时间间隔后再统一发出去还能并发处理多个bulk请求。我使用的参数大致如下每批最多1000条每批数据量不超过5MB每2秒强制flush一次并发请求数2示例配置代码BulkProcessor bulkProcessor BulkProcessor.builder( (request, bulkListener) - esClient.bulkAsync(request, bulkListener), listener ) .setBulkActions(1000) .setBulkSize(new ByteSizeValue(5, ByteSizeUnit.MB)) .setFlushInterval(Duration.ofSeconds(2)) .setConcurrentRequests(2) .build();并发请求数不是越大越好。因为ES写入时需要refresh并发太高反而会导致线程阻塞和内存堆积。个人场景下2并发已经能把吞吐量打到很可观的量级。如果同步的是几十万条数据几千条每秒是可以做到的。这里还要注意一点BulkProcessor的Listener里一定要处理失败请求。只记录“执行成功”但没有检查response中是否有失败detail是大忌。ES bulk接口是允许部分失败的比如某一条数据mapping冲突其他99条正常整体响应里会带着失败原因如果你不解析这一条就静默丢了。3.5 幂等设计与冲突处理同步工具重复执行是常态。定时增量跑完一次可能因为任务卡顿又手动触发了一次上面的状态还没有更新这时候就会重复写入一批数据。如果ES文档的_id是随机生成的那重复写入就等于产生重复文档。所以从一开始就要保证幂等把业务表的主键映射成ES文档的_id。这样无论重复同步多少遍最终索引里每个文档仍然只有一条后写入的会覆盖先写入的。增量同步时我倾向使用update或doc_as_upsert而不是全量的index覆盖。因为增量数据往往只关注变更字段如果用index覆盖某些字段为空时可能把原来的值冲掉。用doc_as_upsert可以只更新给定字段其余字段保持不变这样更安全。IndexRequest request new IndexRequest(index) .id(String.valueOf(row.get(idColumn))) .source(Map.of(title, title, content, content), XContentType.JSON);如果业务上有多个表需要join后再同步到ES可以在SQL里直接join好再按主表主键设置为_id。这样ES文档结构已经反规范化查询时不需要再多次关联。4. 实操过程把一个博客数据同步到ES4.1 环境准备我以最常用的MySQL文章表为例。假设有一张article表字段包括id、title、content、category_id、update_time现在要同步到ES的articles索引。本地环境需要先装好Elasticsearch并启动服务。建议先手动创建索引和mapping不要让ES自动推断否则同步过程中很容易出现类型冲突。我自己用了一个最简的mapping{ mappings: { properties: { title: {type: text}, content: {type: text}, category_id: {type: keyword}, update_time: {type: date} } } }title和content用text类型方便全文搜索category_id用keyword用于精确过滤和聚合update_time用date类型支持按时间过滤。4.2 核心代码实现我用Spring Boot搭建了一个最简工程主要组件有配置读取、任务调度、同步执行器、ES客户端。同步执行器的主流程如下Component public class SyncExecutor { Autowired private JdbcTemplate jdbcTemplate; Autowired private ElasticsearchClient esClient; public void execute(SyncTask task, SyncMode mode) { if (mode SyncMode.FULL) { fullSync(task); } else { incrSync(task); } } private void fullSync(SyncTask task) { long lastId 0L; int totalCount 0; while (true) { ListMapString, Object rows jdbcTemplate.queryForList( SELECT * FROM article WHERE id ? ORDER BY id ASC LIMIT ?, lastId, task.getBatchSize() ); if (rows.isEmpty()) { break; } writeToEs(task.getEsIndex(), rows); totalCount rows.size(); lastId ((Number) rows.get(rows.size() - 1).get(task.getIdColumn())).longValue(); } log.info(全量同步完成任务 {}, 总数 {}, task.getName(), totalCount); } }写入ES的部分使用BulkProcessor统一处理。为了简单演示我在这里直接用同步bulk写入实际项目中用异步BulkProcessor效果更好。4.3 运行效果与日志观察按照上面的结构启动程序后日志大概会像这样输出[main] INFO c.example.sync.SyncExecutor - 开始全量同步任务sync_blog_articles [main] INFO c.example.sync.SyncExecutor - 已同步第 1 批写入 1000 条耗时 310ms [main] INFO c.example.sync.SyncExecutor - 已同步第 2 批写入 1000 条耗时 280ms [main] INFO c.example.sync.SyncExecutor - 全量同步完成总计 50000 条耗时 28s这时可以用ES的Dev Tools查一下索引文档数和数据库的count做对比GET articles/_count如果两边数字一致说明全量同步是通的。之后配置一个定时任务每10分钟跑一次增量同步观察状态表sync_task_state里的last_sync_time是否持续更新。4.4 定时增量同步配置在Spring Boot里定时任务只需要在配置类上加EnableScheduling然后在具体方法上使用Scheduled(cron 0 */10 * * * ?)。每个同步任务可以分别配置cron表达式。Scheduled(cron 0 */10 * * * ?) public void runArticleIncrSync() { SyncTask task taskConfigLoader.load(sync_blog_articles); syncExecutor.execute(task, SyncMode.INCR); }个人项目我一般把增量的cron设置成每5到10分钟一次。如果对实时性要求更高可以缩短到1分钟一轮但要考虑数据库连接消耗。实际上对于个人博客来说10分钟级别的同步已经足够。5. 踩坑记录与问题排查技巧5.1 同步慢且内存抖动第一次把全量同步跑在几十万条数据上时我发现JVM内存曲线像锯齿一样来回波动同步速度也上不来。排查下来主要有三个问题分页查询用了很深offset导致数据库扫描慢一批拿的数据太多内存峰值高ES默认refresh间隔太频繁导致大量小段合并。解决办法是改成分页查询使用keyset方式批次大小控制在500到1000条之间另外在创建索引时把refresh_interval调大到30秒全量同步完再恢复默认值。如果索引数据只用于离线搜索同步期间还可以把副本数设为0同步完成后再改回1写入速度会明显提升。5.2 字段类型不匹配导致写入失败我踩过最经典的一个坑是报错mapper [content] of different type, current_type [text], merged_type [keyword]原因是第一次写入ES时某一条数据的content字段是一个空字符串ES动态映射推断成keyword。后面真正写入长文本时类型又变成text两边冲突就直接写入失败。遇到这种情况最直接的做法是删除索引在同步前手动创建mapping把每个字段的类型都定死不要依赖动态映射。这一点很重要尤其是个人工具想快速建索引时图省事让ES自动推断后续往往会给你埋一堆坑。5.3 增量同步漏数据有一阵子我发现ES里的数据会比MySQL少几条反复核对后定位到是增量字段的时间精度问题。MySQL的datetime精度是秒级同一秒内可能有好几条记录更新而同步任务记录的时间只到秒边界上很容易漏。解决办法有两个第一在记录lastSyncTime时往前偏移几秒让下一轮多查一个时间窗口第二配合自增id作为辅助游标比如查update_time ? OR (update_time ? AND id ?)用复合条件精确控制边界。对我个人项目来说第一种偏移方案已经够用代价是每轮可能重复查少量数据但配合ES的_id幂等不会有副作用。5.4 BulkProcessor静默失败我在第一次跑增量任务时日志显示bulk请求都返回了200但过几天发现ES里某些文档没有更新。后来才知道bulk请求的HTTP状态码只代表请求整体成功内部可能还有逐条失败。必须实现BulkProcessor的Listener在afterBulk里调用response.hasFailures()做检查如果有失败就记录失败message否则问题真的很隐蔽。new BulkProcessor.Listener() { Override public void afterBulk(long executionId, BulkRequest request, BulkResponse response) { if (response.hasFailures()) { log.error(bulk写入存在失败: {}, response.buildFailureMessage()); } } Override public void afterBulk(long executionId, BulkRequest request, Throwable failure) { log.error(bulk请求执行异常, failure); } }5.5 任务状态表的事务边界增量同步是“读状态 - 同步数据 - 写状态”的过程不是强事务。如果同步过程中程序崩溃状态表没有更新下次启动会重复同步一遍。这种设计恰好配合ES的幂等写入重复同步不会产生重复数据。但要注意状态表本身不能出现脏写。如果两个定时任务同时触发了同一个同步任务状态可能被旧任务覆盖。我的做法是在任务启动时先获取一个分布式锁或本地进程内锁个人工具用ReentrantLock就够了保证同一个任务不会并发执行。5.6 扩展从时间轮询到CDC前面提到的方案都是基于时间戳偏移的轮询简单但不够实时。如果你对数据同步的及时性要求更高比如订单支付状态需要秒级反映到ES那就要考虑Canal、Debezium这类基于binlog的CDC方案。在个人项目里接入CDC收益和成本不一定成正比。如果你已经熟悉binlog原理可以尝试把Canal采集到的数据发到本地文件或内存队列再由小工具消费并写入ES。这种方式能捕获删除操作实时性也比轮询好很多。但如果你只是个人博客或小型内容站定时轮询加全量恢复已经足够稳定没必要给系统引入新的故障点。我自己在实际使用中最大的体会是数据同步工具最重要的不是“技术多先进”而是“出问题之后能不能快速恢复”。把状态表、日志、幂等写入这三件事做好一个不起眼的小工具也能变得非常可靠。后续我打算给工具加一个简单的Web控制台可以在页面上查看每个任务最后一次执行时间、手动触发全量或增量同步甚至直接调整增量时间窗口。如果你也在做类似的同步场景建议从最小可用的版本开始先跑通一条链路再把更多数据源加进来一定会省掉很多麻烦。