一、先搞懂核心逻辑:Solr+MapReduce+HDFS的适配思路
很多做数据搜索的朋友可能会有个困惑:平时Solr自己跑索引挺快,但数据量一到百万级、千万级,单机根本扛不住——要么建索引慢到等一天,要么机器内存爆了直接崩。这时候就想到了Hadoop生态里的MapReduce(简称MR),还有能存海量数据的HDFS,那怎么把这仨凑一块,让索引建得又快又稳?
其实核心逻辑很简单:原来Solr是自己读数据、加工、建索引,现在改成让MR帮它分活干。MR天生就是干“把大任务拆成小任务并行跑”的活,比如有100G的待索引数据,MR能拆成10个小任务,每个任务只处理10G,这样10台机器同时跑,速度直接翻10倍。等MR把每个小任务的索引建完,再把这些小索引存到HDFS里,以后Solr要查的时候,直接从HDFS拿索引就行。
举个生活化的例子:就像要给1000本图书编目录,原来你一个人干,得一本本翻、写目录,一天才编10本;现在找10个朋友一起,每人分100本,编完后把每个人的目录都放到一个大仓库(HDFS)里,以后要找某本书的目录,直接去仓库拿就行,这就是MR帮Solr并行建索引的本质。
二、MapReduce下Solr索引并行的完整流程(带可跑示例)
先给大家说清楚整个流程的每一步,再给一个能直接运行的示例,避免大家只懂理论不会落地。
2.1 流程拆解:从数据到并行索引的每一步
第一步:准备待索引的原始数据,比如存在HDFS上的用户行为日志、商品数据,格式可以是CSV、JSON或者自定义的文本格式; 第二步:写MR的Mapper和Reducer代码,Mapper的作用是“把原始数据转成Solr能认的格式”,Reducer的作用是“把Mapper转好的数据,按Solr的规则建小索引”; 第三步:把MR任务提交到Hadoop集群跑,每个Mapper处理一部分数据,每个Reducer生成一个独立的小索引; 第四步:把所有Reducer生成的小索引,上传到HDFS的指定目录,供Solr后续查询加载。
2.2 完整示例(技术栈:Hadoop 3.3.4 + Solr 8.11.2 + Java 8)
这个示例是给HDFS上的商品数据建索引,商品数据是CSV格式,每一行是一个商品,字段包括“商品ID、商品名称、价格、分类”,我们要把这些数据转成Solr的索引。
2.2.1 准备原始数据(先把数据传到HDFS)
先在本地建一个商品数据文件goods.csv,内容如下(模拟100条商品数据,这里只写前5条示意):
1,纯棉T恤,99,服饰
2,无线耳机,199,数码
3,不锈钢保温杯,49,家居
4,运动跑鞋,299,服饰
5,笔记本电脑,4999,数码
然后把这个文件传到HDFS的/input/goods目录:
# 先创建HDFS目录
hdfs dfs -mkdir -p /input/goods
# 上传本地的goods.csv到HDFS
hdfs dfs -put ./goods.csv /input/goods/
2.2.2 写MR的Mapper和Reducer代码
Mapper的作用是把CSV行转成Solr的Document格式(Solr能认的索引数据格式),Reducer的作用是把这些Document按Solr的规则生成小索引。
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Mapper;
import org.apache.solr.common.SolrInputDocument;
import java.io.IOException;
// Mapper类:输入是HDFS上的一行数据,输出是Solr的Document
public class GoodsMapper extends Mapper<LongWritable, Text, Text, SolrInputDocument> {
private Text outKey = new Text("goods_index"); // 所有数据用同一个key,方便Reducer聚合
@Override
protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {
// 把CSV行按逗号拆分,这里假设没有逗号在字段里(实际项目可以用CSV解析工具)
String[] fields = value.toString().split(",");
if (fields.length != 4) {
return; // 跳过格式错误的行
}
// 构造Solr的输入文档(对应商品的四个字段)
SolrInputDocument doc = new SolrInputDocument();
doc.addField("id", fields[0]); // 商品ID,Solr的唯一键
doc.addField("name", fields[1]); // 商品名称
doc.addField("price", fields[2]); // 价格
doc.addField("category", fields[3]); // 分类
// 输出到Reducer,key是goods_index,value是构造好的Solr文档
context.write(outKey, doc);
}
}
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Reducer;
import org.apache.solr.client.solrj.SolrClient;
import org.apache.solr.client.solrj.impl.HttpSolrClient;
import org.apache.solr.common.SolrInputDocument;
import java.io.IOException;
// Reducer类:输入是Mapper传过来的Solr文档,输出是HDFS上的小索引
public class GoodsReducer extends Reducer<Text, SolrInputDocument, Text, Text> {
private SolrClient solrClient;
private String coreName = "goods_core"; // 对应Solr的商品核心(提前在Solr里创建好)
@Override
protected void setup(Context context) throws IOException, InterruptedException {
// 初始化Solr客户端,这里的地址是Solr集群的地址
solrClient = new HttpSolrClient.Builder("http://solr-node1:8983/solr/" + coreName).build();
}
@Override
protected void reduce(Text key, Iterable<SolrInputDocument> values, Context context) throws IOException, InterruptedException {
// 把当前Reducer收到的所有Solr文档批量提交到Solr建索引
for (SolrInputDocument doc : values) {
solrClient.add(doc);
}
// 提交索引(Solr会把这些文档生成索引文件)
solrClient.commit();
// 输出日志,方便查看进度
context.write(new Text("index_built"), new Text("success"));
}
@Override
protected void cleanup(Context context) throws IOException, InterruptedException {
// 关闭Solr客户端,释放资源
if (solrClient != null) {
solrClient.close();
}
}
}
2.2.3 写MR的主类(提交任务到Hadoop)
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;
public class GoodsIndexJob {
public static void main(String[] args) throws Exception {
// 检查输入输出参数(第一个参数是HDFS输入路径,第二个是输出路径)
if (args.length != 2) {
System.err.println("用法: hadoop jar goods-index.jar GoodsIndexJob <输入路径> <输出路径>");
System.exit(1);
}
Configuration conf = new Configuration();
Job job = Job.getInstance(conf, "goods-index-job"); // 任务名称
job.setJarByClass(GoodsIndexJob.class); // 指定主类
// 设置Mapper和Reducer
job.setMapperClass(GoodsMapper.class);
job.setReducerClass(GoodsReducer.class);
// 设置输出的key和value类型
job.setOutputKeyClass(Text.class);
job.setOutputValueClass(Text.class);
// 设置输入输出路径
FileInputFormat.addInputPath(job, new Path(args[0]));
FileOutputFormat.setOutputPath(job, new Path(args[1]));
// 提交任务,等待完成
System.exit(job.waitForCompletion(true) ? 0 : 1);
}
}
2.2.4 打包并提交MR任务
先把上面三个类打包成jar包(比如叫goods-index.jar),然后提交到Hadoop集群跑:
# 提交MR任务,输入路径是HDFS的商品数据,输出路径是HDFS的索引结果目录
hadoop jar goods-index.jar GoodsIndexJob /input/goods /output/goods_index
跑成功后,你会在HDFS的/output/goods_index目录下看到MR生成的结果,同时Solr的goods_core里已经有了新的索引。
三、索引存HDFS后,扩展性和容错带来的新问题
很多人以为把Solr和Hadoop生态搭好就万事大吉了,其实索引存HDFS后,会因为HDFS的扩展性(能存海量数据)和容错(副本机制)带来很多新的麻烦,我给大家拆解清楚:
3.1 扩展性带来的问题:索引越跑越慢,查询卡成狗
HDFS的扩展性是指“能存PB级甚至EB级的数据”,但Solr的索引是有结构的,不是随便存一堆文件就行,具体问题有两个: 第一个是“小索引太多导致的元数据爆炸”。比如你每天跑一次MR建索引,每次拆成100个小索引,跑100天就有10000个小索引,每个小索引都有自己的元数据(比如这个索引对应哪些商品、创建时间、大小),Solr要查的时候,得先扫所有小索引的元数据,再去HDFS拿对应的数据,元数据太多的话,扫的时间比查的时间还长。 举个例子:你去图书馆找一本书,原来只有10个大书架,每个书架的目录贴在门口,扫10次就能找到;现在有10000个小书架,每个书架的目录都要扫一遍,找书的时间直接翻1000倍。
第二个是“HDFS的读取延迟比本地磁盘高”。Solr原来的索引存在本地磁盘,读取速度是毫秒级;现在索引存在HDFS,要先找副本位置、再拉取数据,延迟至少翻2-3倍,要是集群网络拥堵,延迟还会更高,用户查的时候就会感觉卡。
3.2 容错带来的问题:索引不一致,查出来的结果乱
HDFS的容错是靠“副本机制”,每个文件存3份,一份主的两份副的,主的坏了就用副的,但Solr的索引是强一致性的,两者的容错逻辑不兼容,会带来两个大问题: 第一个是“索引更新时的副本冲突”。比如你要更新一个商品的价格,先改主副本的索引,再改两个副副本,要是改到一半集群断网了,主副本改了,副副本没改,Solr查的时候,可能拿到主副本的新价格,也可能拿到副副本的旧价格,结果就乱了。 举个例子:你改了商品的价格是100元,主副本改好了,副副本还没改,用户查的时候,有的用户看到100元,有的看到原来的99元,体验就崩了。
第二个是“MR任务失败导致的索引不完整”。MR任务是并行跑的,要是某个Reducer跑失败了,Solr只会把成功的Reducer生成的索引存到HDFS,失败的那个没生成,这时候Solr查的时候,就会漏掉一部分数据,比如有1000个商品,只建了900个的索引,查的时候就找不到那100个。
四、应用场景、优缺点和注意事项
4.1 适合的应用场景
不是所有情况都适合这么搭,只有满足下面两个条件才划算: 第一个是“数据量特别大”,比如每天新增百万级以上的商品、日志、用户行为数据,单机Solr根本扛不住; 第二个是“对实时性要求不高”,比如电商的商品搜索(每天更新一次索引就行)、日志的历史搜索(查昨天之前的日志),要是要求“改完数据1分钟内就能查到”,这么搭就不合适。
4.2 技术优缺点
优点很明显:一是能处理海量数据,索引速度比单机快很多;二是索引存在HDFS,不会因为机器坏了丢数据;三是成本低,不用买高端的存储设备,用普通服务器就能搭。 缺点也很突出:一是查询延迟高,比单机Solr慢很多;二是索引管理复杂,要处理小索引、副本冲突的问题;三是维护成本高,要同时维护Hadoop集群、Solr集群,出问题的时候要排查两边的日志。
4.3 注意事项
第一个是“提前规划索引的合并策略”,比如每周把所有小索引合并成一个大索引,减少元数据的数量; 第二个是“配置合适的HDFS副本数”,不是副本越多越好,副本多了会占存储空间,一般配置2份副本就够; 第三个是“监控MR任务的运行状态”,要是某个Reducer跑失败了,要及时重跑,避免索引不完整; 第四个是“Solr集群和HDFS集群要同机房”,减少网络延迟,要是跨机房,查询延迟会高到没法用。
五、总结
Solr和Hadoop生态的结合,本质是用MR的并行能力解决大索引的建索引慢的问题,用HDFS解决海量索引的存储问题,但同时也带来了扩展性和容错的新问题。大家在搭之前,一定要先想清楚自己的业务需求:要是数据大、实时性要求不高,这么搭能解决大问题;要是数据小、实时性要求高,还是用单机Solr或者SolrCloud更合适。
最后给大家提个醒:不要为了用技术而用技术,适合自己业务的才是最好的。
评论
围绕“Solr与Hadoop生态集成时,索引构建任务如何在MapReduce框架下并行执行,把索引写入HDFS后,扩展性与容错能力会带来哪些新的问题”参与讨论