简介:这份资源是基于Hadoop的智能购书系统完整项目源码,面向具备Java基础、正在学习大数据处理与推荐算法的开发者及课程设计学生,帮助理解如何用分布式框架搭建一个具备个性化推荐能力的购书平台。压缩包共55个文件,以32个class编译文件与15个java源码为主,另含6个运行日志、1个project工程配置和1个classpath依赖描述,整体约144KB,结构紧凑,便于直接导入IDE阅读与调试。项目围绕HDFS分布式存储与MapReduce并行计算展开,涉及用户行为日志、商品信息与交易记录的处理,并可能结合协同过滤等算法构建推荐引擎,同时引入HBase、Hive等组件完成实时存储与查询分析。已有190人学习关注,适合作为大数据入门到进阶的实践参考,可据此梳理Hadoop生态各组件的协作方式、推荐逻辑的实现思路以及Java编写MapReduce作业的工程组织方法。
1. 从零搭一套基于 Hadoop 的智能购书系统:它到底解决什么问题
很多人第一次听到「基于 Hadoop 的智能购书系统」,脑子里浮现的是个电商网站,其实重点根本不在前端页面。它真正要解决的是:当购书平台的用户行为日志、订单流水、图书元数据涨到单机数据库扛不住的时候,怎么用 Hadoop 这套分布式存储和计算框架,把「猜你喜欢什么书」这件事算出来。换句话说,这是一个把 Hadoop 离线计算能力套在图书零售场景上的课程设计级项目,也是很多高校大数据专业最典型的综合实践题。
它适合三类人:正在做 Hadoop 课程设计、需要一套能跑通全流程参考方案的学生;想从单机 MySQL 转分布式、拿一个完整场景练手的后端或数据开发;以及面试前想补一段「HDFS + MapReduce + 推荐逻辑」实战经历的求职者。这篇文章不讲空泛概念,从伪分布式搭建一路讲到推荐结果落库,把参数、命令和踩过的坑都摊开说。
2. 环境先立住:Hadoop 伪分布式搭建与 Docker 镜像两条路
2.1 为什么课程设计优先选伪分布式而不是全分布式
全分布式集群(多台机器、NameNode 和 DataNode 分离)听起来更「高级」,但对一个购书系统课程设计来说,投入产出比很低。伪分布式是在一台机器上把 NameNode、DataNode、ResourceManager、NodeManager 全部跑起来,逻辑上等价于一个最小集群,能完整验证 HDFS 读写和 MapReduce 作业提交。你真正要证明的是「我懂分布式计算流程」,而不是「我有三台服务器」。
我一般会建议:本地虚拟机或云主机给 4 核 8G 起步,磁盘留 50G。JDK 用 8(Hadoop 3.x 对 JDK 8 兼容最稳),Hadoop 选 3.3.x 系列。下面这套流程在 Ubuntu 22.04 上验证过。
先做基础准备,创建专用用户并配好 SSH 免密,这是 Hadoop 启动脚本依赖的:
# 创建 hadoop 用户并授权 sudo useradd -m -s /bin/bash hadoop sudo passwd hadoop sudo usermod -aG sudo hadoop # 切换到 hadoop 用户,配置本机免密登录 su - hadoop ssh-keygen -t rsa -P '' -f ~/.ssh/id_rsa cat ~/.ssh/id_rsa.pub >> ~/.ssh/authorized_keys chmod 600 ~/.ssh/authorized_keys ssh localhost # 能免密进去说明配置成功逻辑说明:Hadoop 的 start-dfs.sh 会通过 SSH 去启动各节点进程,即使是伪分布式也走这套机制,所以免密是硬前提。参数上-P ''表示空密码,-t rsa指定密钥类型,别用默认交互式一路回车,脚本化更省事。
接着装 JDK 和 Hadoop,配环境变量:
# 解压到 /opt 并改属主 sudo tar -zxvf jdk-8u381-linux-x64.tar.gz -C /opt/ sudo tar -zxvf hadoop-3.3.6.tar.gz -C /opt/ sudo chown -R hadoop:hadoop /opt/jdk1.8.0_381 /opt/hadoop-3.3.6 # 编辑 ~/.bashrc,追加以下内容 export JAVA_HOME=/opt/jdk1.8.0_381 export HADOOP_HOME=/opt/hadoop-3.3.6 export PATH=$PATH:$JAVA_HOME/bin:$HADOOP_HOME/bin:$HADOOP_HOME/sbin export HDFS_NAMENODE_USER=hadoop export HDFS_DATANODE_USER=hadoop export HDFS_SECONDARYNAMENODE_USER=hadoop export YARN_RESOURCEMANAGER_USER=hadoop export YARN_NODEMANAGER_USER=hadoop逻辑说明:Hadoop 3.x 的启动脚本会检查这些*_USER变量,不配会直接报「please define HDFS_NAMENODE_USER」之类的错,这是新手最常见的翻车点。改完执行source ~/.bashrc生效。
2.2 四个核心配置文件的参数怎么填
Hadoop 伪分布式要改的配置文件都在$HADOOP_HOME/etc/hadoop/下,一共四个。参数填错是最容易导致 DataNode 起不来的原因,逐个说清楚。
core-site.xml指定默认文件系统和临时目录:
<configuration> <property> <name>fs.defaultFS</name> <value>hdfs://localhost:9000</value> </property> <property> <name>hadoop.tmp.dir</name> <value>/opt/hadoop-3.3.6/tmp</value> </property> </configuration>hdfs-site.xml设置副本数为 1(伪分布式只有一块盘,设 3 会一直报副本不足):
<configuration> <property> <name>dfs.replication</name> <value>1</value> </property> <property> <name>dfs.namenode.name.dir</name> <value>/opt/hadoop-3.3.6/data/namenode</value> </property> <property> <name>dfs.datanode.data.dir</name> <value>/opt/hadoop-3.3.6/data/datanode</value> </property> </configuration>mapred-site.xml让 MapReduce 跑在 YARN 上,yarn-site.xml配 ResourceManager 主机和 NodeManager 辅助服务:
<!-- mapred-site.xml --> <configuration> <property> <name>mapreduce.framework.name</name> <value>yarn</value> </property> </configuration> <!-- yarn-site.xml --> <configuration> <property> <name>yarn.nodemanager.aux-services</name> <value>mapreduce_shuffle</value> </property> <property> <name>yarn.resourcemanager.hostname</name> <value>localhost</value> </property> </configuration>参数说明:dfs.replication=1是伪分布式必须的;mapreduce.framework.name=yarn决定作业提交到 YARN 而不是本地;yarn.nodemanager.aux-services=mapreduce_shuffle是 Shuffle 阶段依赖的服务,漏配会导致 Reduce 卡在 0%。
2.3 格式化与启动,以及 Docker 镜像这条捷径
首次启动前必须格式化 NameNode,且只能格式化一次:
hdfs namenode -format start-dfs.sh start-yarn.sh jps # 应看到 NameNode/DataNode/SecondaryNameNode/ResourceManager/NodeManagerjps是判断集群健康的第一道关。如果 DataNode 没出现,九成是hadoop.tmp.dir里残留了旧的 clusterID,删掉 tmp 和 data 目录重新格式化即可。
如果你不想折腾虚拟机,用 Docker 镜像更快。常见做法是拉一个带 Hadoop 的镜像,把配置目录挂载出来:
docker run -d --name hadoop-pseudo \ -p 9870:9870 -p 8088:8088 -p 9000:9000 \ -v $(pwd)/hadoop-conf:/opt/hadoop/etc/hadoop \ hadoop:3.3.6逻辑说明:9870 是 HDFS Web UI,8088 是 YARN Web UI,9000 是 RPC 端口。挂载配置目录的好处是改参数不用进容器。注意容器内也要保证 SSH 和*_USER变量配好,否则一样起不来。
3. 数据从哪来、怎么进 HDFS:购书日志的采集与建模
3.1 智能购书系统的数据模型设计
一个能算推荐的购书系统,最少需要三张核心数据:用户行为日志、订单明细、图书元数据。行为日志是推荐算法的燃料,字段设计直接决定后面能算什么。
我一般会设计成这样的行为表结构,用制表符或逗号分隔存成文本文件:
| 字段 | 含义 | 示例 |
|---|---|---|
| user_id | 用户编号 | U10023 |
| book_id | 图书编号 | B2045 |
| behavior_type | 行为类型 | view / cart / buy / rate |
| behavior_time | 行为时间戳 | 1716883200 |
| duration | 停留秒数 | 45 |
图书元数据表包含 book_id、title、category、author、price;订单表包含 order_id、user_id、book_id、amount、order_time。这三张表在 HDFS 上以目录区分,比如/bookstore/logs/、/bookstore/books/、/bookstore/orders/。
提示:行为类型用英文枚举而不是中文,能避免后续 MapReduce 里编码不一致导致的乱码,这是血泪经验。
3.2 用 put 和 distcp 把数据送进 HDFS
本地造好测试数据后,上传到 HDFS:
# 建目录 hdfs dfs -mkdir -p /bookstore/logs /bookstore/books /bookstore/orders # 上传本地文件 hdfs dfs -put ./data/user_behavior.txt /bookstore/logs/ hdfs dfs -put ./data/books.txt /bookstore/books/ hdfs dfs -put ./data/orders.txt /bookstore/orders/ # 验证 hdfs dfs -ls /bookstore/logs/ hdfs dfs -cat /bookstore/logs/user_behavior.txt | head -5逻辑说明:-put适合小批量测试数据;如果数据量大到几个 G,用-put会慢,常见做法是先用hdfs dfs -mkdir建目录,再用distcp从另一个 HDFS 或对象存储并行拷贝。参数上-p保留权限和时间戳,-f覆盖已存在文件。
数据进 HDFS 后,一个绕不开的概念是 InputSplit。MapReduce 作业提交时,框架会把大文件切成一个个 InputSplit,每个 Split 交给一个 Map 任务处理。默认 Split 大小等于一个 HDFS Block(128M),所以一个 1G 的日志文件会被切成 8 个 Split,起 8 个 Map。理解这点很重要:如果你的数据只有几 MB,却起了几十个 Map,说明小文件太多,每个小文件单独成一个 Split,Map 任务启动开销远大于计算本身。解决办法是先用一个合并作业把小文件合成大文件,或者调大mapreduce.input.fileinputformat.split.minsize。
3.3 用 MapReduce 统计图书热度
先写一个最基础的热度统计作业,统计每本书被浏览的次数,作为推荐的输入之一:
public class BookHotStat { // Mapper:输出 <book_id, 1> public static class HotMapper extends Mapper<LongWritable, Text, Text, IntWritable> { private final static IntWritable one = new IntWritable(1); private Text bookId = new Text(); @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { // 日志格式:user_id \t book_id \t behavior_type \t time \t duration String[] fields = value.toString().split("\t"); if (fields.length >= 3 && "view".equals(fields[2])) { bookId.set(fields[1]); context.write(bookId, one); } } } // Reducer:累加同一本书的浏览次数 public static class HotReducer extends Reducer<Text, IntWritable, Text, IntWritable> { @Override protected void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException { int sum = 0; for (IntWritable val : values) { sum += val.get(); } context.write(key, new IntWritable(sum)); } } public static void main(String[] args) throws Exception { Configuration conf = new Configuration(); Job job = Job.getInstance(conf, "book hot stat"); job.setJarByClass(BookHotStat.class); job.setMapperClass(HotMapper.class); job.setReducerClass(HotReducer.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(IntWritable.class); FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[1])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }逻辑说明:Mapper 里做了行为类型过滤,只统计 view,避免把购买行为混进热度。split("\t")的分隔符必须和上传数据一致,用错分隔符会导致 fields 长度不对,Map 输出为空,最后结果文件是空的——这是最隐蔽的坑之一。Reducer 的Iterable是框架自动按 key 分组后的迭代器,不用自己排序。
打包提交:
# 编译打包 javac -classpath $(hadoop classpath) -d classes BookHotStat.java jar -cvf bookhot.jar -C classes/ . # 提交作业 hadoop jar bookhot.jar BookHotStat /bookstore/logs /bookstore/output/hot # 查看结果 hdfs dfs -cat /bookstore/output/hot/part-r-00000 | head参数说明:$(hadoop classpath)自动带上 Hadoop 所有依赖 jar,省得手动一个个加。输出目录必须不存在,否则作业直接报FileAlreadyExistsException,这是新手反复踩的坑。
4. 推荐算法落地:从共现矩阵到图书推荐结果
4.1 基于物品的协同过滤为什么适合购书场景
推荐算法有很多,购书场景我优先选基于物品的协同过滤(ItemCF)。原因很实际:图书的数量相对稳定,用户数量却可能很大且不断变化,ItemCF 计算的是「书和书之间的相似度」,这个矩阵可以离线算好、定期更新,不用每次用户来了重算。而 UserCF 在用户量爆炸时相似度矩阵会大到算不动。
ItemCF 的核心逻辑是:如果很多用户同时买了 A 和 B,那 A 和 B 就相似;给买了 A 的用户推荐 B。落到 MapReduce 上,分两步:第一步统计物品共现,第二步算相似度并生成推荐。
4.2 用 MapReduce 算物品共现矩阵
先按用户分组,把每个用户买过的书两两组合:
// Mapper:以 user_id 为 key,book_id 为 value public static class CoMapper extends Mapper<LongWritable, Text, Text, Text> { @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String[] fields = value.toString().split("\t"); if (fields.length >= 3 && "buy".equals(fields[2])) { context.write(new Text(fields[0]), new Text(fields[1])); // <user, book> } } } // Reducer:同一用户的书两两配对输出 <bookA:bookB, 1> public static class CoReducer extends Reducer<Text, Text, Text, IntWritable> { @Override protected void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException { List<String> books = new ArrayList<>(); for (Text val : values) { books.add(val.toString()); } // 两两组合 for (int i = 0; i < books.size(); i++) { for (int j = 0; j < books.size(); j++) { if (i == j) continue; context.write(new Text(books.get(i) + ":" + books.get(j)), new IntWritable(1)); } } } }逻辑说明:Reducer 里把同一用户的所有书收集到 List 再两两配对,输出<书A:书B, 1>。这里有个性能隐患:如果一个用户买了几百本书,两两组合是 O(n²),会撑爆内存。常见做法是加一个上限,比如只取最近购买的 50 本,或者用mapreduce.reduce.memory.mb调大内存并配合-Xmx。
第二步再跑一个作业,把相同<书A:书B>的计数累加,得到共现次数,然后除以各自热度开方得到余弦相似度。相似度公式:
sim(A,B) = cooc(A,B) / sqrt(count(A) * count(B))其中 cooc 是共现次数,count 是单本书的购买次数(前面热度统计的结果可以直接复用)。
4.3 生成推荐列表并落库
有了相似度矩阵,给每个用户推荐时,遍历他买过的书,找出每本书最相似的 TopN,加权累加得分,排除已买的,取分数最高的若干本:
// 伪代码逻辑,实际用 MapReduce 或 Spark 实现 for (String book : userBoughtBooks) { List<Pair> similar = similarityMap.get(book); // 该书最相似的 N 本 for (Pair p : similar) { if (userBoughtBooks.contains(p.bookId)) continue; // 排除已买 scoreMap.put(p.bookId, scoreMap.getOrDefault(p.bookId, 0.0) + p.sim); } } // 按 score 排序取 Top10 作为推荐结果最终推荐结果写成user_id, book_id, score三列,用hdfs dfs -get拉回本地,再通过 JDBC 批量写入 MySQL 供前端查询。批量插入时用rewriteBatchedStatements=true参数能显著提速,这是很多人不知道的 MySQL 连接串技巧。
注意:推荐结果落库前一定要去重和过滤冷门书,否则新用户会收到一堆没人买过的书,体验很差。
5. 避坑与排查:这套系统最容易翻车的五个地方
5.1 DataNode 启动后立刻消失
现象:start-dfs.sh后jps里没有 DataNode,或者启动几秒后进程没了。原因:hadoop.tmp.dir或dfs.datanode.data.dir里残留了上一次格式化的 clusterID,和当前 NameNode 的 clusterID 不一致,DataNode 拒绝加入。解决:停掉集群,删掉 tmp 和 data 目录,重新hdfs namenode -format,再启动。记住格式化只能做一次,重复格式化会让已有数据全部失效。
5.2 作业卡在 Map 或 Reduce 100% 不动
现象:YARN Web UI 上作业一直停在 map 100% reduce 0%,或者 map 阶段长时间不动。原因:常见有三种——yarn.nodemanager.aux-services没配mapreduce_shuffle,Reduce 拿不到 Map 输出;数据倾斜,某个 key 的数据量远超其他,单个 Reduce 拖慢整体;内存不足导致容器被 kill 后重试。解决:先查yarn-site.xml的 shuffle 配置;数据倾斜可以在 Mapper 里给 key 加随机前缀打散,Reduce 后再去掉;内存问题调mapreduce.map.memory.mb和mapreduce.reduce.memory.mb。
5.3 输出目录已存在导致作业直接失败
现象:提交作业报org.apache.hadoop.mapred.FileAlreadyExistsException: Output directory ... already exists。原因:MapReduce 不允许输出目录预先存在,防止覆盖已有结果。解决:每次提交前删掉输出目录hdfs dfs -rm -r /bookstore/output/hot,或者在代码里用FileSystem.exists判断后删除。别图省事直接改源码跳过检查,会埋下数据覆盖的隐患。
5.4 中文图书标题乱码
现象:结果文件里图书标题显示成问号或方块。原因:源文件编码是 GBK,而 Hadoop 默认按 UTF-8 读取。解决:上传前用iconv -f GBK -t UTF-8 books.txt > books_utf8.txt转码,或者在建表时就统一用 UTF-8。这个坑在课程设计答辩现场翻车率极高,因为本地编辑器默认编码各不相同。
5.5 小文件过多拖垮 NameNode
现象:HDFS 上几万个几 KB 的小文件,NameNode 内存飙升,作业启动几十个 Map 却每个只处理几行。原因:每个小文件在 NameNode 占约 150 字节元数据,且每个文件单独成一个 InputSplit。解决:用hadoop archive打成 HAR 包,或者跑一个合并作业把小文件合成大文件,再或者在上游采集时就做批量写入。购书系统的日志如果按天按小时切得太碎,这个问题几乎必然出现。
6. 进阶技巧:用 Combiner 和分区优化让作业快一倍
前面跑通流程只是及格线,真正让这套智能购书系统在生产数据量下还能跑得动,靠的是两个优化点:Combiner 和自定义 Partitioner。这两个东西不改变业务逻辑,但能把网络传输和磁盘 IO 砍掉一大半。
Combiner 本质是「Map 端的 Reduce」,在 Map 输出落盘前先做一次局部聚合。以热度统计为例,同一个 Map 任务里可能处理了同一本书的上千条浏览记录,如果不加 Combiner,这上千个<book, 1>全都要通过网络传给 Reduce;加了 Combiner 后,先在本地累加成<book, 1000>再传,网络传输量直接降三个数量级。用法很简单,在驱动类里加一行:
job.setCombinerClass(HotReducer.class);逻辑说明:Combiner 的输入输出类型必须和 Reducer 一致,所以直接复用 HotReducer 就行。但要注意,Combiner 只适用于满足结合律和交换律的操作,求和、求最大值可以,求平均值不行——因为平均值不能简单地对局部平均值再平均。这是很多人想当然用错的地方。
Partitioner 决定 Map 输出交给哪个 Reduce。默认是HashPartitioner,按 key 的 hash 取模。如果你的数据本身倾斜严重,比如某本畅销书的记录占了 30%,默认分区会让一个 Reduce 累死、其他 Reduce 闲着。解决办法是自定义 Partitioner,把热门 key 打散:
public static class SkewPartitioner extends Partitioner<Text, IntWritable> { @Override public int getPartition(Text key, IntWritable value, int numPartitions) { // 对热门书加盐,分散到不同 Reduce if (isHotBook(key.toString())) { return (key.hashCode() & Integer.MAX_VALUE) % numPartitions; } return (key.hashCode() & Integer.MAX_VALUE) % numPartitions; } }逻辑说明:& Integer.MAX_VALUE是为了把 hashCode 的负数转成正数,否则取模会得到负分区号直接报错。这个位运算技巧在 Hadoop 源码里到处都是,记住它比记公式管用。
验证优化效果,最直接的办法是看 YARN 作业的 Counter。提交作业后打开 8088 端口的 Web UI,找到你的 application,看Reduce shuffle bytes这个指标——加 Combiner 前后对比,通常能降 60% 到 90%。另一个指标是Map output records和Reduce input records的比值,比值越大说明 Combiner 效果越好。
我自己的习惯是:任何 MapReduce 作业上线前,先不加任何优化跑一遍记录基线,再加 Combiner 跑一遍,最后加 Partitioner 跑一遍,三次的 Counter 截图存下来。这样既知道优化有没有效果,答辩或复盘时也有据可查。这套基于 Hadoop 的智能购书系统,真正值钱的不是推荐算法多花哨,而是你能把分布式计算的每个环节调明白、把每个坑填上。希望帮到你。
本文还有配套的精品资源,点击获取