2021年MongoDB与Apache Solr数据同步及索引最优方案咨询
MongoDB 同步到 Apache Solr 实用方案
已弃用 DIH 的替代方案
目前官方推荐的稳定同步方案主要有三类,适配不同的使用场景:
- 全量批量导入:
mongoexport+ Solr Bulk API 组合,无额外依赖,操作简单 - 自定义同步逻辑:SolrJ 客户端自主开发,适配复杂的字段映射、清洗需求
- 实时增量同步:MongoDB Change Stream + SolrJ 组合,实现毫秒级数据同步
可运行示例
1. 单集合全量导入
适合一次性同步单个Mongo集合的全量数据:
- 导出Mongo集合为JSON Lines格式(默认每行一个JSON文档,符合Solr导入要求)
mongoexport --uri="mongodb://你的Mongo地址:27017/目标库名" --collection=目标集合名 --type=json --out=./mongo_data.jsonl
- 直接导入Solr,提前确认Solr Core的schema字段和Mongo字段对应,可通过
jq工具做字段重命名/过滤
# 示例将Mongo的_id字段映射为Solr要求的唯一id字段 cat ./mongo_data.jsonl | jq '{id: ._id, title: .title, content: .content, create_time: .create_time}' | curl -X POST -H "Content-Type: application/json" "http://你的Solr地址:8983/solr/目标Core名/update?commitWithin=10000" --data-binary @-
commitWithin=10000表示10秒内自动提交,比每次导入强制提交性能高3倍以上,适合大批量数据导入
2. 多集合批量导入
同一个库下多个集合批量同步的Shell脚本示例,自动处理不同集合同_id的冲突问题:
#!/bin/bash MONGO_URI="mongodb://你的Mongo地址:27017/目标库名" SOLR_URL="http://你的Solr地址:8983/solr/目标Core名/update?commitWithin=10000" # 填入要同步的集合列表 COLLECTIONS=("集合1" "集合2" "集合3") for col in "${COLLECTIONS[@]}" do echo "正在同步集合:$col" mongoexport --uri="$MONGO_URI" --collection="$col" --type=json | \ # 增加source字段标记集合来源,id拼接集合名避免冲突 jq --arg col "$col" '. + {source: $col, id: "\($col)_\(._id)"}' | \ curl -X POST -H "Content-Type: application/json" "$SOLR_URL" --data-binary @- done
3. 实时增量同步
生产环境常用的全量+增量组合方案,全量导入完成后启动Change Stream监听Mongo数据变更,实时同步到Solr,核心Java代码示例:
import com.mongodb.client.MongoClients; import com.mongodb.client.MongoClient; import com.mongodb.client.MongoDatabase; import com.mongodb.client.model.changestream.ChangeStreamDocument; import org.apache.solr.client.solrj.impl.Http2SolrClient; import org.apache.solr.common.SolrInputDocument; import org.bson.Document; public class MongoSolrSyncService { public static void main(String[] args) { MongoClient mongoClient = MongoClients.create("mongodb://你的Mongo地址:27017"); MongoDatabase db = mongoClient.getDatabase("目标库名"); Http2SolrClient solrClient = new Http2SolrClient.Builder("http://你的Solr地址:8983/solr/目标Core名").build(); // 监听全库所有集合的变更,也可指定单个集合监听 db.watch().forEach(change -> { if (change.getOperationType() == null) return; String collName = change.getNamespace().getCollectionName(); String docId = collName + "_" + change.getDocumentKey().getObjectId("_id").toString(); try { switch (change.getOperationType()) { case INSERT, UPDATE -> { Document fullDoc = change.getFullDocument(); SolrInputDocument doc = new SolrInputDocument(); doc.addField("id", docId); doc.addField("source", collName); doc.addField("title", fullDoc.getString("title")); doc.addField("content", fullDoc.getString("content")); // 其余字段按需映射 solrClient.add(doc); } case DELETE -> solrClient.deleteById(docId); } solrClient.commit(); } catch (Exception e) { e.printStackTrace(); } }); } }
注意事项
- 字段映射规则复杂的场景,可以单独配置映射文件,不用硬编码在代码/脚本中
- 超大集合全量导入可以加
--batchSize参数分批导出,避免内存溢出 - 生产环境建议增加同步失败重试、断点续传逻辑,避免网络波动导致数据丢失
内容的提问来源于stack exchange,提问作者scofx
相关产品推荐
相关产品推荐

