ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Hadoop MapReduce实战:电影网站用户性别预测教案

Hadoop MapReduce实战:电影网站用户性别预测教案 简介针对Hadoop大数据开发基础课程项目案例教案以“电影网站用户性别预测”为主线面向大数据技术类专业师生帮助学习者将KNN算法与MapReduce分布式计算结合完成从数据预处理、分类器构建到模型评价的完整流程。内容系统涵盖KNN算法原理与实现步骤、MapReduce编程实现KNN、准确率/召回率/F1-score等评价指标并设计了引导性、探究性与拓展性问题如MapReduce连接两份文件、K值选取等以及理论教学和实验教学的分阶段安排适合教师备课或学生自学参考。资源为单份PDF文件共1个文件压缩包整体约24KB内容精炼但结构完整包含教学目标、重点难点、教学过程设计、教材与参考资料等模块。已有1078人浏览学习。通过学习可快速理清MapReduce KNN的实现逻辑掌握连接用户数据/评分数据/电影数据、清洗缺失值与异常值、划分训练集验证集测试集、模型寻优等关键环节并了解单机版KNN与分布式KNN的差异可直接用于课堂讲授或课程设计参考。1. 用 Hadoop 跑通电影网站用户性别预测这个教案项目到底在教什么电影网站要按性别做推荐偏偏用户注册时性别字段大量缺失。这个教案级的项目案例就是根据用户已有的评分行为和历史属性用 Hadoop 上的 MapReduce 把缺失的性别预测出来。数据用经典的 MovieLens 三张表就可以开跑算法主流选朴素贝叶斯训练过程就是频次统计天然能拆成 Map 和 Reduce 两段比在 Hadoop 上硬跑神经网络划算得多。适合刚学完 HDFS 命令、想上手第一个完整大数据开发作业的人也适合带实训课的老师直接拿来做任务拆解。真做一遍你会发现模型代码不是门槛跨表聚合和特征清洗才是真正花时间的地方没跑通一次完整作业你对作业提交到 YARN 的流程理解基本停留在面试题层面。2. 把性别预测拆成 MapReduce 能算的步骤评分特征、分类模型与数据流设计2.1 为什么拿性别预测当大数据开发入门项目从教学角度这是目前我能想到的最合适的二分类案例。第一任务足够简单性别只有 M 和 F 两种取值预测错了可以用准确率直接量化不需要解释复杂的业务指标。第二输入数据天然分表用户信息、电影信息、评分记录分散在三个文件里做特征就得先跨表聚合这正是 MapReduce 最擅长、也是企业里最常干的事。第三这套流程能覆盖 HDFS 操作、Mapper/Reducer 编写、作业调参、结果回传四个环节一个案例串起整个大数据开发基础。很多人一上来就想着用逻辑回归甚至深度学习这在课堂上容易翻车梯度下降需要多轮迭代每一轮都要跑一遍完整的 MapReduce 作业伪分布式下单次 shuffle 的耗时就会让新手失去耐心。朴素贝叶斯不一样它是“一步式”训练——统计完频次概率就出来了。后文所有的设计都围绕这一个选择展开。2.2 三张表的字段口径与标签定义MovieLens 是电影评分研究里常用的公开数据教学版一般包含 users、movies、ratings 三个文件。我一般先把它们整理成统一的 CSV字段口径如下文件字段在本项目中的用法users用户ID、性别、年龄、职业、邮编性别做标签年龄和职业做特征邮编丢弃movies电影ID、标题、类型类型用来给评分数据挂特征ratings用户ID、电影ID、评分、时间戳按用户聚合成行为特征需要注意三张表之间存在外键关系ratings 里的用户ID和电影ID分别关联 users 和 movies。做特征聚合时我们先把 ratings 与 movies 关联得到“用户对某类型电影的评分明细”再按用户分组算出行为特征而不是直接把三张表做笛卡尔积那样数据量会膨胀到无法收敛。2.3 特征怎么选离散字段分段行为字段聚合性别预测的可用特征分两类。第一类是用户属性直接取 users 表里的字段年龄和职业。年龄是连续值但在朴素贝叶斯里最好分段比如 18、18-24、25-34、35-44、45-54、55 六段每段变成一个离散取值。职业本身就是分类编号可以直接用。邮编在入门级别通常丢弃因为类别太多、每个类别的样本太少统计出来的概率不稳定。第二类是行为特征来自 ratings 表按用户聚合后的统计量。我常用的有四个评分次数、平均评分、评分方差、以及各类型电影观看占比。评分次数反映活跃度平均评分反映口味倾向类型占比是预测性别最有区分度的特征——比如动作片占比高的用户男性概率明显更大。聚合逻辑本质上是按用户ID做分组统计这正是 MapReduce 里“合并去重”思路的延伸同一个用户的多条评分在 Map 阶段以用户ID为 key 归并在 Reduce 阶段算出均值、方差和占比。2.4 朴素贝叶斯的 MapReduce 实现逻辑朴素贝叶斯的出发点很简单给定一组特征计算用户属于男性或女性的后验概率取大者作为预测结果。公式是 P(性别|特征) ∝ P(性别) × Π P(特征_i|性别)。要拿到右边的一串概率训练阶段只需要做两类统计统计每个性别的用户数得到先验概率 P(性别)。统计每个性别下各特征取值的出现次数得到条件概率 P(特征_i|性别)。这两类统计在 MapReduce 里是同一个模式Mapper 读一条用户样本把“性别”和“性别:特征取值”作为 key 输出Value 固定为 1Reducer 对相同 key 的计数累加最后把累加值除以对应性别的总样本数就得到概率。整个过程不需要迭代一个 MapReduce 作业就能完成训练。计算 P(性别|特征) 时有一个工程细节多个条件概率连乘几十个特征之后结果会小到浮点数下溢。所以实际操作中不会直接乘而是取对数把乘法变成加法log P(性别) Σ log P(特征_i|性别)。这也是后面预测代码里要处理的关键点。2.5 两个作业串起来的整体数据流整个项目我习惯拆成两个 MapReduce 作业。作业一拿 users 表训练输出先验概率和条件概率写到 HDFS 上的一个文本文件作业二读同一个 users 表模拟线上特征已知、标签缺失的场景每个 Mapper 的 setup 阶段把模型文件加载到内存然后逐条预测。之所以强调两个作业是因为这样职责清晰训练作业的结果是模型文件预测作业的结果是 user_id 与预测性别两边的输出都能单独检查。数据流可以描述为原始三张表 → 清洗成 CSV → 上传 HDFS → 作业一统计概率 → 作业二加载模型并预测 → 结果拉回本地评估。后面的章节就按这个顺序展开。实际做的时候你会发现作业一跑完基本就成功了一大半预测作业的坑集中在模型文件的路径和格式上。3. Hadoop 伪分布式环境与 MovieLens 数据准备从环境搭建到 HDFS 可读3.1 伪分布式搭建的最小配置三个配置文件和两条启动命令教学和单机开发场景我一般建议用 Hadoop 的伪分布式模式一台机器上同时跑 NameNode、DataNode、ResourceManager、NodeManager够验证完整作业流程又不需要三台虚拟机。准备步骤没有什么玄学先装 JDK配好 JAVA_HOME把 Hadoop 解压到一个不含空格的路径比如 /opt/hadoop然后改三个文件。core-site.xml 决定默认文件系统!-- core-site.xml默认走 HDFS并指定 NameNode 的 RPC 地址 -- configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property property namehadoop.tmp.dir/name value/data/hadoop_tmp/value /property /configuration这里 fs.defaultFS 写 localhost 而不是主机名是为了避免后续从 Windows 客户端访问时 DNS 解析失败hadoop.tmp.dir 是 NameNode 和 DataNode 数据落盘的位置首次使用前必须确保这个目录存在并且当前用户可写。hdfs-site.xml 里最关键是副本数伪分布式只有一个 DataNode副本数设成 3 会导致大量块一直处于等待复制状态所以改成 1。yarn-site.xml 保持默认即可如果要限制单个作业内存才会去调 yarn.nodemanager.resource.memory-mb。配置完成后第一次启动必须格式化 NameNode# 首次使用必须执行之后不要随意重复 bin/hdfs namenode -format # 启动 HDFS 和 YARN sbin/start-dfs.sh sbin/start-yarn.sh # 用 jps 确认五个进程都在 jpsjps 输出里应该能看到 NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager 五个进程。少了哪个都不能继续往下走这是第一次提交作业前必做的自检。格式化这个动作相当于系统初始化重复格式化会清理掉已有数据属于“后悔药”操作改配置前想清楚自己的数据是否还有用。3.2 数据预处理从 Windows 编辑到 Linux 运行的编码坑教学场景里很多人习惯在 Windows 上用 IDEA 写代码数据文件也是 Windows 上编辑的这就引入第一个翻车点MovieLens 原始文件是 GBK 编码或带 BOM 的 UTF-8直接丢给 Hadoop 跑出来的字段带 \ufeff 前缀join 对不上。所以数据必须先过一遍清洗脚本把编码统一、字段统一、脏行滤掉。# coding: utf-8 # 清洗三个原始文件统一编码、过滤异常行、转成逗号分隔的 CSV def clean(src, dst, sep::, field_num3): with open(src, rb) as fin, open(dst, w, encodingutf-8) as fout: for raw in fin: line raw.decode(utf-8, errorsignore).strip() if not line: continue fields line.split(sep) if len(fields) ! field_num: continue # 字段数不对的直接丢不让脏数据进 HDFS fout.write(,.join(fields) \n) if __name__ __main__: clean(users.dat, users_clean.csv, field_num5) clean(ratings.dat, ratings_clean.csv, field_num4)这段脚本的特点是先把每一行按 :: 切分然后用字段数做硬校验users 必须是 5 列ratings 必须是 4 列错一行就跳过。参数 field_num 是每类文件的字段数约束比只做 strip 过滤可靠得多。出来的 CSV 统一是 UTF-8 无 BOM后面的 Java 代码里直接用 split(,) 解析不会踩编码坑。3.3 把清洗结果放上 HDFS目录规划与文件检查数据清洗完先别急着写代码先把目录结构和文件放到位。HDFS 上我习惯按作业名建目录并在下面分 input/output 两层hdfs dfs -mkdir -p /movie/input/users hdfs dfs -mkdir -p /movie/input/ratings hdfs dfs -mkdir -p /movie/input/movies hdfs dfs -put users_clean.csv /movie/input/users/ hdfs dfs -put ratings_clean.csv /movie/input/ratings/ hdfs dfs -put movies_clean.csv /movie/input/movies/ # 检查块分布确认文件真的进了 HDFS hdfs fsck /movie/input/ratings -files -blockshdfs fsck 这条命令值得多说一句它输出每个文件被切成了多少个块、每个块在哪个 DataNode。教学场景经常有人问“明明 put 成功了为什么程序读不到”多半是路径写错或者 put 到了本地文件系统——先用 fsck 确认路径存在比反复改代码省时间。另外注意HDFS 的基础操作ls、put、get、cat 这几条命令在这一步里会全部用到。3.4 第一次提交作业前的自检进程、日志与常见失败信号环境搭好后先别急着提交自己的 MR 程序跑一个官方自带的示例作业验证整个链路。示例在 Hadoop 安装目录的 share/hadoop/mapreduce 下常见做法是用 wordcount 探路hadoop jar share/hadoop/mapreduce/hadoop-mapreduce-examples-*.jar \ wordcount /movie/input/users /movie/output/demo这个作业的作用不是数单词而是验证从客户端到 ResourceManager 的调度链路是否正常。如果这个作业能跑完说明环境没问题后面自己的代码出错就是程序问题如果连示例都卡住就要去看 YARN 的日志位置在 logs/userlogs/ 下按 job 编号分目录里面会写清楚是内存不够、端口被占还是节点未注册。第一次独立跑的时候把“看日志”变成条件反射能少走一大半弯路。4. 训练和预测的 MapReduce 代码Mapper、Reducer 与作业提交参数4.1 训练作业一个 Mapper 同时产出先验和条件概率训练作业的输入是 users_clean.csv每行五列user_id, gender, age, occupation, zip。我设计的 key 分两类一类是 prior:性别用于统计男女总数一类是 性别:特征取值用于统计条件概率。这样只用一组 Mapper 和 Reducer就能同时得到先验和条件概率输出文件写在一起预测时再分开解析。public class GenderTrain { public static class TrainMapper extends MapperLongWritable, Text, Text, IntWritable { private final static IntWritable ONE new IntWritable(1); private final Text outKey new Text(); Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String[] f value.toString().split(,); if (f.length 3) return; // 防御性检查脏数据跳过 String gender f[1].trim(); // M 或 F这就是标签 String ageGroup ageToGroup(f[2]); // 年龄分段 // 先验概率计数器prior:M / prior:F outKey.set(prior: gender); context.write(outKey, ONE); // 条件概率计数器M:A / F:B 这类 key outKey.set(gender : ageGroup); context.write(outKey, ONE); } private String ageToGroup(String age) { int a Integer.parseInt(age.trim()); if (a 18) return A; if (a 25) return B; if (a 35) return C; if (a 45) return D; if (a 55) return E; return F; } } public static class TrainReducer extends ReducerText, IntWritable, Text, Text { Override protected void reduce(Text key, IterableIntWritable values, Context context) throws IOException, InterruptedException { long sum 0; for (IntWritable v : values) sum v.get(); context.write(key, new Text(Long.toString(sum))); } } }Mapper 里输出的 value 统一是 1Reducer 只做累加这是频次统计最朴素的写法。两个细节值得注意一是 outKey 复用同一个 Text 对象避免每行 new 两个对象造成的 GC 压力二是 ageToGroup 把年龄映射成 A 到 F 的枚举值而不是直接用数字这样分类特征不会出现“35 和 36 被当成两个完全不同的类别”的问题。训练作业的输出就是预测作业需要的模型文件。4.2 预测作业在 setup 阶段把模型加载进每个 Mapper预测作业的输入可以复用 users_clean.csv但语义上要假设性别字段不可用。每个 Mapper 在处理数据前先通过 -files 参数加载模型文件把里面的每一行解析成概率。这里用到 MapReduce 最常见的“外部参数分发”机制DistributedCache。下面的代码里模型文件通过 -files ...#model.txt 分发任务工作目录里会有一个叫 model.txt 的软链接所以直接用 FileReader 就能读。public class GenderPredict { public static class PredictMapper extends MapperLongWritable, Text, Text, Text { private final MapString, Double cond new HashMap(); private double totalM 0.0, totalF 0.0; private double priorM, priorF; private final Text outVal new Text(); private final Text outKey new Text(result); Override protected void setup(Context context) throws IOException, InterruptedException { // model.txt 来自 -files 参数里的 # 别名 try (BufferedReader br new BufferedReader(new FileReader(model.txt))) { String line; while ((line br.readLine()) ! null) { String[] kv line.split(\\t); if (kv[0].startsWith(prior:M)) totalM Double.parseDouble(kv[1]); else if (kv[0].startsWith(prior:F)) totalF Double.parseDouble(kv[1]); else cond.put(kv[0], Double.parseDouble(kv[1])); // 条件频次 } } // 先验转成对数后面预测时用加法替代乘法 priorM Math.log(totalM / (totalM totalF)); priorF Math.log(totalF / (totalM totalF)); } Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String[] f value.toString().split(,); if (f.length 3) return; String ageGroup ageToGroup(f[2]); // 拉普拉斯平滑频次1分母加上年龄分段数 6 double scoreM priorM Math.log((cond.getOrDefault(M: ageGroup, 0.0) 1) / (totalM 6)); double scoreF priorF Math.log((cond.getOrDefault(F: ageGroup, 0.0) 1) / (totalF 6)); String pred scoreM scoreF ? M : F; outVal.set(f[0] , pred); context.write(outKey, outVal); } } }预测代码里有两个容易写错的地方。第一setup 中要区分“先验计数”和“条件频次”先验计数用于算男女比例条件频次则需要除以对应性别的总样本数并且做平滑。第二scoreM 和 scoreF 全程用对数累加就是为了避免几十个条件概率连乘后下溢成 0。输出 key 统一写成 result让所有预测结果集中到一个 Reducer输出文件就是一份完整的 user_id,predicted_gender 结果表。4.3 作业提交参数hadoop jar、-files 与输出目录规则代码写完打成 jar 包用 hadoop jar 提交。我把两个作业的提交命令一起说明# 训练作业从 CSV 统计男女先验和各特征条件概率 hadoop jar gender.jar com.example.GenderTrain \ -D mapreduce.job.reduces1 \ /movie/input/users /movie/output/model # 预测作业把模型文件通过 -files 分发到每个 Mapper hadoop jar gender.jar com.example.GenderPredict \ -D mapreduce.job.reduces1 \ -files /movie/output/model/part-r-00000#model.txt \ /movie/input/users /movie/output/predict每个参数都有讲究整理成一张表方便查阅参数作用常见误用-D mapreduce.job.reduces指定 Reducer 个数训练/预测固定 1 即可不写则按集群默认值可能分区紊乱-files 路径#别名把 HDFS 上的模型文件分发到每个任务节点只写路径不带别名程序里找不到文件名输入路径Mapper 读取的 HDFS 目录可写多个写成本地路径作业直接失败输出路径Reducer 结果的落地目录必须不存在目录已存在会报 FileAlreadyExistsExceptionshuffle 阶段对新手来说像个黑匣子但从提交参数看它就是把 Map 输出按 key 分区、排序、拷贝给 Reduce 的过程。输出目录必须不存在这是 Hadoop 刻意设计的防覆盖机制想重复跑作业得先删掉上一次的输出目录。4.4 查看结果从 HDFS 拉回本地验证预测作业跑完后输出目录下会有 part-r-00000 文件。我一般用 getmerge 把多个分区文件合并拉回本地而不是一个个 gethdfs dfs -getmerge /movie/output/predict ./predict_result.csv wc -l predict_result.csvgetmerge 会把目录下所有 part 文件按字典序拼接成单个文件。在伪分布式下只有一个 Reducer拉回来的是一份完整结果在完全分布式下如果有多个 Reduce 输出getmerge 的价值才真正体现。拉回来之后立刻看前 10 行确认输出确实是 user_id,M/F 的格式再进入评估环节。5. 性别预测作业的 5 个高频坑数据倾斜、内存溢出与路径权限排查5.1 卡在 Map 100% 但 Reduce 一直 0%现象作业的 Map 阶段已经到 100%Reduce 进度却一直是 0%日志里没有异常就是不动。新手通常以为是死机其实多半卡在 shuffle。原因Map 输出到 Reduce 输入之间有个数据拷贝过程如果 Map 输出的数据量比较大或集群只有一个节点且网络回环很慢Reduce 要等所有 Map 的输出准备好才开始。另一种常见原因是某个 Mapper 失败后一直在重试。解决先看 YARN 的 userlogs确认没有任务失败然后调大 reduce 拉取并发参数例如在作业提交时加 -D mapreduce.reduce.shuffle.parallelcopies8再把 Mapper 的输出做压缩减小 shuffle 数据量。对这个项目把特征压缩到个位数级别shuffle 的耗时是可接受的。5.2 模型概率全是负无穷或 NaN现象训练作业跑完模型文件里全是 -Infinity 或 NaN预测结果也全是同一个性别。原因某个性别下的某个年龄段样本数为 0取对数时变成 log(0) 即负无穷。比如样本里女性 35-44 岁段一条都没有公式里 P(女|35-44) 就成了 0对数直接炸掉。解决一是清洗数据时过滤掉非法的年龄值二是做拉普拉斯平滑统计频次时每个类别都加 1公式从 count / total 改成 (count 1) / (total 类别数)。平滑量虽小但对空类别兜底很关键。这是所有概率模型项目都会遇见的经典问题教案里提前写好这一条能省不少排错时间。5.3 数据倾斜几个 Reducer 里有一个特别慢现象设置了多个 Reducer其中一个跑了很久其他的早就结束整个作业被拖住。原因性别本身就不均衡有的数据集男女比例 7:3加上动作片占比这类特征在男性里高度集中某个 key 的累计频次远大于其他 key对应 Reducer 就处理不动。解决先看倾斜的 key 是什么。如果只是训练阶段的频次统计最简单的方案是把 Reducer 数量设小或直接设为 1因为频次统计需要全局归并如果是预测阶段输出结果可以加一个前缀随机数做“加盐”把热点 key 拆散下一轮作业再按真实 key 归并。加盐会让结果分到多个输出文件后续用 getmerge 接回来就行。5.4 Connection refused 到 9000 端口现象在本机提交作业时报连接 localhost:9000 被拒绝但 Hadoop 进程明明已经用 jps 查过了。原因core-site.xml 里 fs.defaultFS 配的是主机名而客户端所在的机器解析不了这个主机名或者 hosts 文件里没有对应映射。Windows 上用 IDEA 提交时尤其常见。解决把 fs.defaultFS 统一写成 hdfs://localhost:9000如果必须用主机名就在 /etc/hostsWindows 是 C:\Windows\System32\drivers\etc\hosts里加一行“本机IP 主机名”。配置文件里的地址要保持一致不然格式化后 DataNode 注册不了还会引出 clusterID 不一致的问题。5.5 重复格式化导致 DataNode 起不来现象改了配置后重新执行 namenode -formatNameNode 起来了DataNode 却一直在日志里报错或直接退出。原因format 会为 NameNode 生成新的 clusterID但 DataNode 的数据目录里还保留着旧 clusterID两边对不上DataNode 拒绝注册。解决格式化之前把 hadoop.tmp.dir 对应的目录整个删掉再重新 format保证 NameNode 和 DataNode 的 clusterID 同步生成。这个操作会清空 HDFS 上的所有数据所以只适合数据还没价值或已经下载到本地的阶段。伪分布式环境下这是最常翻车的点没有之一属于写进血泪经验的那类问题。6. 用正确率和混淆矩阵验证模型再谈训练集/测试集划分这个后悔药模型跑完不能只看“好像挺准”。我把预测结果拉回本地后会立刻写一段 Python把预测性别和真实性别对比输出准确率和混淆矩阵。这段脚本很短但必须做否则你永远不知道模型是“真的会”还是“只会猜男性”——类别不平衡时全猜 M 的准确率也可能有 70%。# 评估预测结果predict_result.csv 只有 user_id,predicted import pandas as pd from sklearn.metrics import confusion_matrix, accuracy_score label pd.read_csv(users_clean.csv, names[uid, gender, age, occ, zip]) pred pd.read_csv(predict_result.csv, names[uid, pred]) data label.merge(pred, onuid) print(accuracy:, accuracy_score(data.gender, data.pred)) print(confusion_matrix(data.gender, data.pred))如果评估时用的是训练同一份数据准确率会虚高因为模型已经见过这些样本。真正可信的做法是划分训练集和测试集而且要按用户ID切不能按评分记录切。同一个用户的多条评分如果同时落在训练和测试两侧等于把答案提前泄露给了模型这是最容易忽略但影响最大的坑也是我强调的“后悔药”开始评估前先把划分想好不然只能回头重新训练。这个项目作为大数据开发教案价值不在模型本身而在完整的技术链路伪分布式环境搭建、数据清洗、HDFS 操作、MapReduce 编程、作业提交、结果评估。把这个链路跑熟再去看 YARN 调度细节、Hive 或 Spark 的实现时你的收获会完全不一样。我自己每次重讲这个案例都会重新写一遍评估脚本提醒自己不看测试集、不先调参再选模型这是一条踩坑换来的习惯。希望帮到你。本文还有配套的精品资源点击获取
返回列表