多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

基于Hadoop与朴素贝叶斯的电影网站用户性别预测

基于Hadoop与朴素贝叶斯的电影网站用户性别预测 简介面向大数据与机器学习初学者的基于HadoopKNN的电影网站用户性别预测完整实现。程序围绕用户性别分类任务提供从数据预处理到模型建模、评价的一站式代码预处理阶段封装为JAR包部署到Hadoop集群即可运行KNN计算阶段支持本地执行便于调试并针对大规模数据给出运行耗时提示。压缩包共78个文件以28个Java源码、30个编译后的class文件及5个可运行JAR包为主辅以CLASSPATH/PROJECT/PROPERTIES配置文件和DAT数据文件整体约5.86MB目录按数据预处理、拆分数据、建模评价等模块划分便于按需取用。当前已有3038人学习下载。通过这份资料读者可以掌握Hadoop环境下KNN算法的工程化落地方法理解分布式数据预处理与本地模型计算如何衔接并可直接修改路径复用性别预测流程适合课程设计、毕业设计或推荐系统入门实践。1. 基于 hadoop 的电影网站用户性别预测先想清楚这个“基于”落在哪看到这个标题很多人第一反应是“用 Hadoop 训练一个性别分类模型”。这类课程设计或模拟项目做多了之后我的判断恰恰相反Hadoop 在这个组合里最合适的角色是数据底座负责存储、清洗和分布式统计真正承担分类的是一个能被 MapReduce 天然拆解的朴素贝叶斯模型。电影网站的用户评分日志动辄百万行单机 Pandas 在聚合阶段就会吃力而性别预测本身又是一个特征极其稀疏的二分类问题用不上神经网络。于是最常见的可靠路线是把用户评分记录整理成宽表上传 HDFS用训练、预测两个 MapReduce Job 跑通全流程。这篇文章按这条路线往下拆数据怎么组织、两个 Job 怎么分工、参数怎么调、哪些坑让作业翻车适合正在做 Hadoop 课程设计或者想把大数据平台和基础分类算法放进同一个程序里的从业者。想明白“基于”落在哪后面每一步都有依据。2. 先把用户评分日志整理成宽表特征列、标签列和 MovieTypes 的拆分规则2.1 四类字段怎么划分不是所有列都能直接进模型电影网站的原始数据一般散落在三张表里用户表存性别电影表存类型和名称评分表记录用户对电影的评分和打分时间。做性别预测时最忌讳的是在 MapReduce 里反复 join 这三张表因为 shuffle 会重复扫描 HDFS作业时间翻倍数据倾斜还特别难排查。我一般会先做一次预处理把三张表拼成一张宽表一行对应一条评分记录性别和电影类型直接冗余在行内。宽表至少包含以下字段字段名示例值在模型中的角色user_idu_1001聚合键一个用户对应多行genderM / F标签列预测目标movie_idm_4422调试用不进入特征movie_typesAction|Adventure|Sci-Fi核心特征多类型用管道符分隔rating4.2辅助特征本方案不使用timestamp1556889600训练/测试切分的时间依据这里有一个课程设计里经常犯的错有人会把 gender 也做成特征模型准确率立刻冲到 95% 以上看起来漂亮但本质上是让模型直接读标签属于典型的数据泄漏。特征列和标签列的边界必须分开gender 只出现在标签位置在任何一行特征里都不能出现。另一个关键点在于 movie_types 是多值列。如果整列作为一个 token 统计Action|Adventure 和 Action|Comedy 会被当成两个完全独立的特征动作片样本被拆得七零八落条件概率估计会非常稀疏。所以预处理的第一步就是决定好“拆分粒度”。2.2 特征选择为什么“类型偏好”比评分均值更值得信任电影类型和用户性别之间的偏好差异在很多数据集里都能观察到比如动作片和科幻片的用户群体里男性占比更高爱情片和剧情片的女性占比更高。这个规律虽然不是绝对但作为朴素贝叶斯的条件概率来源是够用的。相比之下评分均值反而不适合做第一特征。原因是评分行为受个人打分习惯影响极大有人习惯只给高分有人只看了烂片也打 4 分直接做连续特征需要离散化离散化边界又是一层玄学。所以在第一个可落地版本里我把评分列留成辅助特征模型只吃类型序列。一个用户看过的所有电影类型拼在一起就构成一个“词袋”性别是这个词袋的类别标签。这个选择还有一层好处MapReduce 对词袋统计天然友好。每一条评分记录可以独立地输出一个gender, type计数对不需要跨行聚合Combiner 也能安全使用。如果将来想把评分加权进去改动也只在 Mapper 里把计数 1 换成 rating两个 Job 的整体结构不用重建。2.3 清洗命令与 HDFS 上传落盘之前先清理三个隐患原始评分文件经常带表头、空值和大小写混乱的性别编码还有类型列里的多余空格。这些不清理干净后面模型文件里会出现一批“看起来一样”的孤儿特征。我习惯先用一个 awk 命令把它洗成规整的 TSV# 1) 过滤表头、空行统一性别编码去掉类型列内的空格 awk -F \t NR1 NF6 $2! { gender($2male||$2M)?M:F; gsub(/ /,,$4); print $1\tgender\t$3\t$4\t$5\t$6 } raw_ratings.tsv clean_ratings.tsv # 2) 确认清洗后的行数和字段格式 wc -l clean_ratings.tsv head -3 clean_ratings.tsv # 3) 建目录并上传到 HDFS hdfs dfs -mkdir -p /input/sex/raw hdfs dfs -put clean_ratings.tsv /input/sex/raw/这个脚本里NF6保证字段数不足的行直接丢弃gsub(/ /,,$4)是为了避免“Action ”和“Action”在后续聚合时成为两个键。性别编码统一成 M/F 则是为了让训练和预测阶段看到完全一致的标签取值。上传时要注意一个常见误解Map 任务数不是人为指定的而是由 InputFormat 决定的。TextInputFormat 默认按文件字节区间生成 InputSplit一个 split 对应一个 Mapsplit 的边界尽量对齐 HDFS 块。比如当前块大小是 128MB一个 300MB 的文件通常会切成 3 个 split。如果上传目录里有几百个小文件每个小文件至少一个 split就会白白启动几百个 Map这就是课程设计里最常见的小文件问题。遇到这种情况可以把小文件合并成大文件再上传或者直接改用 CombineFileInputFormat这个话题在第 4 章参数部分展开。数据落盘只是第一步真正决定这个项目能不能讲清楚的是第 3 章的建模拆解。3. 训练与预测两个 MapReduce Job朴素贝叶斯的分布式拆解3.1 模型文件格式与平滑参数先行定义朴素贝叶斯在这个场景下的公式不复杂。设类别 C 取值为 M 或 F用户看过的类型集合是 t1, t2, …在条件独立假设下预测分数是类别的对数先验与每个类型对数条件概率之和score(C) log P(C) Σ log P(tj | C)哪个类别的分数高就预测哪个性别。训练阶段要统计的只有两类量先验 P(C)也就是男、女用户各自占比条件概率 P(t | C)也就是某个性别的用户看过某类电影的比例。条件概率直接用频次相除会有一个隐患某个类型在训练集里没出现过条件概率算出来是 0预测时一乘整个分数就变负无穷。所以必须做拉普拉斯平滑P(t | C) (count(t, C) α) / (count_total(C) α * V)这里的 α 默认取 1V 是去重后的电影类型总数。α 不要写死在代码里我一般把它作为 Driver 的命令行参数传入做交叉验证时需要频繁调这一项。平滑后的概率全部取对数保存在模型文件里这样预测 Mapper 读取后可以直接做累加不用切换对数空间。模型文件格式自己定义就行我用的格式是两行前缀区分先验和条件概率PRIOR M -0.712 PRIOR F -0.619 COND M Action -3.281 COND F Action -4.012这里所有数值都是 log 空间Job2 的 Mapper 读到后直接放进 HashMap。3.2 Job1 统计计数Mapper 拆键、Reducer 汇总Job1 的输入就是第 2 章生成的宽表输出是性别, 类型和性别, 总量的计数。Mapper 的核心逻辑是把一行记录拆成多个计数键public class TypeCountMapper extends MapperLongWritable, Text, Text, LongWritable { private final Text outKey new Text(); private final LongWritable one new LongWritable(1L); Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String line value.toString(); if (line.startsWith(user_id)) { return; // 跳过宽表表头 } String[] fields line.split(\t, -1); if (fields.length 6 || fields[1].isEmpty()) { return; // 字段缺失的行直接丢弃 } String gender fields[1]; for (String type : fields[3].split(\\|)) { type type.trim(); if (type.isEmpty()) { continue; // 空类型不参与统计 } outKey.set(gender \t type); context.write(outKey, one); // (性别, 类型) 计数 outKey.set(gender \t*TOTAL*); context.write(outKey, one); // (性别, 总量) 计数 } } }这段代码里split(\t, -1)的第二个参数是刻意保留末尾空字段防止行尾字段缺失时数组长度被 Java 静默截断。每行输出两路键一路是具体的性别, 类型一路是性别,TOTAL后者用于计算条件概率的分母。TOTAL哨兵是我习惯用的写法比起用 Hadoop Counter 更直观因为在 Reducer 里可以直接看到全量计数后续 buildModel 读取时也统一走文件扫描。Reducer 只做求和public class TypeCountReducer extends ReducerText, LongWritable, Text, LongWritable { Override protected void reduce(Text key, IterableLongWritable values, Context context) throws IOException, InterruptedException { long sum 0L; for (LongWritable value : values) { sum value.get(); } context.write(key, new LongWritable(sum)); } }求和语义下 Combiner 可以安全复用同一个 Reducer 类直接job1.setCombinerClass(TypeCountReducer.class)就行能明显减少 shuffle 到 Reducer 的数据量。有一点要提醒Combiner 只在满足交换律和结合律的操作里安全计数没问题求均值不行所以 Job2 里不能想当然地套用这个做法。Job1 跑完之后输出目录里是纯频次还需要一个 buildModel 步骤把频次转成模型文件。这个步骤在 Driver 里完成读取所有 part 文件计算平滑概率后写出static void buildModel(Path countDir, Path modelPath, Configuration conf, double alpha) throws IOException { MapString, Long condCount new HashMap(); MapString, Long genderTotal new HashMap(); FileSystem fs FileSystem.get(conf); for (FileStatus part : fs.listStatus(countDir)) { String name part.getPath().getName(); if (!name.startsWith(part)) continue; // 逐行解析 (性别\t类型\t计数)*TOTAL* 行放入 genderTotal // 其余行放入 condCountkey 为 性别 \t 类型 } SetString vocab new HashSet(); for (String k : condCount.keySet()) { vocab.add(k.split(\t, -1)[1]); } try (PrintWriter out new PrintWriter(fs.create(modelPath))) { long totalUsers genderTotal.getOrDefault(M, 0L) genderTotal.getOrDefault(F, 0L); for (String g : new String[]{M, F}) { double prior (double) genderTotal.getOrDefault(g, 0L) / totalUsers; out.println(PRIOR\t g \t Math.log(prior)); long total genderTotal.getOrDefault(g, 0L); for (String type : vocab) { long cnt condCount.getOrDefault(g \t type, 0L); double prob (cnt alpha) / (total alpha * vocab.size()); out.println(COND\t g \t type \t Math.log(prob)); } } } }这段逻辑里V 直接取去重后的类型集合大小比从 Counter 读更省事。如果训练数据特别大单机扫描模型文件可能吃内存但常见课程设计数据规模下没有压力。3.3 Job2 预测setup 读模型、map 算分数、reduce 定性别Job2 的输入是测试集宽表模型文件通过 DistributedCache 分发给每个 Map 容器。Mapper 在 setup 阶段把模型文件读进两个 Map然后在 map 方法里对每条记录累加分数public class SexPredictMapper extends MapperLongWritable, Text, Text, Text { private final MapString, Double logPrior new HashMap(); private final MapString, MapString, Double logCond new HashMap(); Override protected void setup(Context context) throws IOException, InterruptedException { Path[] localFiles context.getLocalCacheFiles(); if (localFiles null) return; for (Path p : localFiles) { try (BufferedReader reader new BufferedReader( new InputStreamReader(FileSystem.getLocal(context.getConfiguration()) .open(new Path(p.toString()))))) { String line; while ((line reader.readLine()) ! null) { String[] parts line.split(\t); if (parts.length 3) continue; if (PRIOR.equals(parts[0])) { logPrior.put(parts[1], Double.parseDouble(parts[2])); } else if (COND.equals(parts[0]) parts.length 4) { logCond.computeIfAbsent(parts[1], k - new HashMap()) .put(parts[2], Double.parseDouble(parts[3])); } } } } } private double lookup(String gender, String type) { return logCond.getOrDefault(gender, Collections.emptyMap()) .getOrDefault(type, context.getConfiguration() .getDouble(sex.default.logprob, -20.0)); } Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { String[] fields value.toString().split(\t, -1); if (fields.length 6 || fields[1].isEmpty()) return; double m logPrior.getOrDefault(M, 0.0); double f logPrior.getOrDefault(F, 0.0); for (String type : fields[3].split(\\|)) { type type.trim(); if (type.isEmpty()) continue; m lookup(M, type); f lookup(F, type); } context.write(new Text(fields[0]), new Text(fields[1] \t m \t f)); } }lookup方法里那个-20.0的默认值是一个调参细节测试集中出现训练集没有的类型时直接返回一个比较小的对数概率而不是抛异常数值上也几乎等价于零概率。这个默认值不用写死在代码里通过sex.default.logprob配置项传入交叉验证时可以观察它对结果的影响。为什么不在 Mapper 端直接比较 m 和 f 输出结果因为同一个用户有多行记录如果在 Mapper 内部做判断需要缓存一个用户的全部历史内存不可控。所以 Mapper 只输出中间分数把累加和比较交给 Reducerpublic class SexPredictReducer extends ReducerText, Text, Text, Text { Override protected void reduce(Text key, IterableText values, Context context) throws IOException, InterruptedException { double m 0.0, f 0.0; String realGender ; for (Text v : values) { String[] parts v.toString().split(\t); if (parts.length 3) continue; realGender parts[0]; m Double.parseDouble(parts[1]); f Double.parseDouble(parts[2]); } String predict m f ? M : F; context.write(key, new Text(predict \t realGender)); } }Reducer 的输入 key 是 user_idvalue 是该用户每一行记录带来的分数增量。累加完成后比较两个类别的总分数输出三列用户 ID、预测性别、真实性别。到这里整个模型已经能跑出一个可评估的结果文件但离“在集群上稳定跑通”还差一段配置和调参的路。4. 作业提交与参数调优本地验证、伪分布式到集群的完整闭环4.1 本地模式快速跑通最小命令拿到代码先不要急着上集群本地模式能省下大量排错时间。把 MapReduce 的 framework 切成本地模式输入输出直接走本地文件系统# 打包Hadoop 依赖用 provided 作用域 mvn clean package -DskipTests # 本地模式运行输入输出都是本地路径 hadoop jar target/sexpredict-1.0.jar \ -D mapreduce.framework.namelocal \ -D fs.defaultFSfile:/// \ /data/clean_ratings.tsv /tmp/sex/model /tmp/sex/out /tmp/sex/count 1.0这个命令里四个路径参数依次是输入宽表、模型输出目录、预测结果目录、Job1 计数目录最后的 1.0 是拉普拉斯平滑系数。本地模式下没有 HDFS 也无所谓程序里所有 FileSystem 操作都走 file:/// 协议可以用来验证 Mapper 和 Reducer 逻辑是否正确也可以快速确认数据格式有没有问题。如果连 Hadoop 环境都还没搭可以起一个伪分布式集群在一台机器上同时启动 NameNode、DataNode 和 ResourceManager再把数据传入 HDFS# 伪分布式启动 HDFS 和 YARN start-dfs.sh start-yarn.sh # 检查输入是否已经上传 hdfs dfs -ls /input/sex/raw伪分布式和本地模式最大的差别在于本地模式根本没有走网络传输和真实的分布式调度很多和序列化、节点通信相关的坑只有伪分布式或真实集群才能暴露出来。我一般会在本地模式跑通后立刻用同一份打包产物在伪分布式上再跑一遍两组结果对比能快速区分“逻辑问题”和“环境问题”。4.2 三个关键参数分片、Reducer 数量与 Combiner跑通之后就要面对性能问题。一个 MapReduce 程序的性能百分之七八十由分片大小、Reducer 数量和 Combiner 决定这三个参数在这个项目里都有明确的选择依据。参数默认行为在这个项目里的设置建议mapreduce.input.fileinputformat.split.maxsize对齐 HDFS 块输入文件规整时不用动小文件多时配合 CombineFileInputFormat 使用mapreduce.job.reduces1Job1 设 1 方便读模型Job2 按数据量设 6~20 个Combiner未设置Job1 复用 Reducer 类做求和Job2 不设先说分片。如果上传目录里是几十个小文件每个小文件一个 splitMap 数量会远大于你想要的并行度。这时候可以调整 split 上限让多个小文件合并进同一个分片或者直接换输入格式job.setInputFormatClass(CombineFileInputFormat.class); CombineFileInputFormat.setMaxInputSplitSize(job, 128L * 1024 * 1024);再说 Reducer 数量。Job1 我建议设成 1这样计数结果只出现在一个 part-r-00000 文件里buildModel 读取逻辑最简单。数据量大到单个 Reducer 吃不下时再拆拆多个的话 buildModel 需要扫描整个 part 目录逻辑也不复杂。Job2 的 Reducer 则是按用户 ID 哈希分区的数量设成 6~20 个比较合理因为每个用户只在同一个分区内累加不会出现跨分区合并的问题。设成用户数那是在自找麻烦一个用户一个 Reducer 会白白启动上万个任务。Combiner 的设置在 3.2 里已经提过Job1 可以复用 TypeCountReducer 类因为频次求和满足交换律和结合律。Job2 不能设 Combiner因为 Reducer 做的是两个分数的比较比较操作不满足结合律设了 Combiner 会直接算出错误结果。4.3 用日志和 Counter 定位问题从空跑到数据倾斜作业跑完不等于跑对。我每次提交完作业第一件事是看两样东西作业结束时客户端打印的 Counter 摘要以及 YARN 上的日志。Counter 能直观看到 map 输入记录数、shuffle 字节数、reduce 输出记录数。如果 map 输入记录数是 0说明 InputFormat 没读到数据先检查输入路径如果 reduce 输出记录数远小于预期用户数说明聚合键设计有问题。查看运行中或已完成作业的日志# 先拿 application ID yarn application -list # 按 ID 拉取全部容器日志 yarn logs -applicationId application_1234567890数据倾斜在这个项目里也很常见。性别样本本身不平衡比如男性用户记录占 80%那么“M\tTOTAL”这个键对应的数据量会明显大于“F\tTOTAL”哈希分区后两个 Reducer 负载差异会拉大。观察现象是某个 Reducer 运行时间远超其他 Reducershuffle 字节数也偏高。如果影响到了可接受范围可以给计数键加盐比如拆成“性别类型随机后缀”的多个键第二轮再聚合但课程设计阶段通常不需要走到这一步。先确认倾斜原因再决定要不要加盐这是一条值得记住的排错顺序。5. 训练预测全流程避坑五个课堂不讲的翻车点5.1 预测结果全输出男性类别不平衡把条件概率压死了现象结果文件里 90% 以上的用户都被预测成 MF 几乎永远不出现。原因通常是训练数据里男性用户占比过高先验 P(M) 远大于 P(F)即使某个用户看过大量女性主导的类型片log 求和后总分仍被先验拖向 M。另一层隐蔽原因是如果按“记录数”而不是“用户数”计算先验男性用户打分频率更高时会进一步放大先验偏差。解决先把判断标准从“比较 m 和 f”改成“比较两者差值是否超过阈值”阈值通过验证集调比如 0.3 的对数分数差才判 M否则判 F同时在报告指标时不要只看准确率用宏平均 F1 或者混淆矩阵评估类别不平衡时准确率没有任何参考价值。5.2 随机切分训练/测试集让准确率虚高同一个人被拆成了两份样本现象随机抽样 90% 做训练、10% 做测试准确率跑到 92%一换到真实场景立刻掉到 70% 左右。原因同一个用户的评分记录同时出现在训练集和测试集模型在训练时见过的类型偏好模式在测试时被同一个人的其他记录“剧透”了属于数据泄漏。解决切分单位必须是 user_id而不是行。常见的做法是给每个用户算一个 hashhash 值落在某个区间内的用户整体进测试集# 给宽表追加一列 fold_id按用户哈希取余 awk -F \t { fold $1 % 10; # user_id 是 u_1001 时先去掉前缀 print $0\tfold } clean_ratings.tsv clean_ratings_fold.tsv如果更贴近真实业务场景还可以按用户注册时间或评分时间切分窗口这部分在第 6 章展开。5.3 Reducer 内存溢出把整个模型文件全量读进每个任务现象Job2 启动后少量 Map 任务反复失败查看日志是 Java heap space。原因模型文件本身没多大但 DistributedCache 分发后Mapper 的 setup 里只做了一次读取问题不在模型文件本身而在于有人顺手把测试集也读进了内存或者在 setup 里用 FileSystem.open 反复打开远程 HDFS 文件高峰期把 NameNode 连接打满。解决setup 里只读模型文件测试数据必须通过 Mapper 的 value 逐行进入读取模型时优先用 context.getLocalCacheFiles()不要每次 map 都去 open 一个远程路径。注意getLocalCacheFiles 在较新的 API 中建议用 job.addCacheFile 搭配 PathFilter 来替代但课程设计环境里两种写法都能工作关键是别在 map 方法里做文件 IO。5.4 本地能跑、上集群就失败依赖冲突和 HADOOP_HOME 双坑现象本地模式输出完全正常打包后提交到集群直接 NoClassDefFoundError或者提示找不到主类。原因有一半是打包时把 Hadoop 的依赖打了进去和集群运行时环境冲突另一半是提交节点的 HADOOP_HOME 或 PATH 没配对hadoop jar 命令用了错误版本的入口。解决Maven 工程里 Hadoop 相关依赖统一标成 provided让作业运行时使用集群自带库提交前检查 hadoop version 和集群版本一致HADOOP_HOME 指向部署目录。这个坑在简历里说出来很加分因为大多数人都是在本地 IDE 里跑通就以为万事大吉。5.5 预测分数永远 TIE类型列分隔符在传输中被改写现象结果文件里大量行的预测性别是 TIE两个分数一模一样。原因测试集的 movie_types 列在预处理或复制过程中被 Tab 键覆盖管道符变成了制表符Mapper 里 split(\|) 永远切不出类型分数始终只有先验加和两边相等。解决提交作业前先用hdfs dfs -cat /input/sex/test/* | head抽查原始行确认第 4 列分隔符还是管道符。这类问题很蠢但出现频率极高我现在的习惯是数据一上传就先抽样看一条完整记录不只看 wc -l。6. 收尾验证把“跑通”变成“可靠”的三个进阶动作一个课程设计或简历项目光有准确率数字是撑不住追问的。面试官大概率会问“你怎么知道这个模型不是过拟合”“为什么不用随机切分”。所以最后建议补三个验证动作成本都不高。第一个做按用户的折交叉验证。把用户按 hash 分成 5 份每次拿一份当测试集循环提交 5 次作业for fold in 0 1 2 3 4; do hadoop jar sexpredict-1.0.jar \ -D sex.fold$fold \ /input/sex/fold_$fold/test /tmp/sex/model_$fold /tmp/sex/out_$fold \ /tmp/sex/count_$fold 1.0 done5 次准确率的均值和方差比单次准确率诚实得多。第二个在 Reducer 里用 Counter 输出混淆矩阵。真实性别和预测性别的组合有四种M-M、M-F、F-M、F-F用四个 Counter 累加作业跑完直接读数字F1 分数一算就有。类别不平衡时这个动作尤其重要它能让你看到模型到底是在瞎猜还是真的有区分度。第三个加一个时间窗口验证。电影类型偏好会随档期变化用前 80% 时间段的评分训练、预测后 20% 时间段的用户更接近真实部署。窗口换几次准确率波动如果超过 10%说明特征里隐含了时间敏感信号需要重新考虑要不要把时间差本身做成特征。我现在做这类预测项目第一个 Job 永远是输出样本分布的 Counter先看类不平衡再看结果写完作业先跑本地模式再上伪分布式最后才提交集群。这套顺序帮我挡掉了 70% 的低级问题。希望帮到你。本文还有配套的精品资源点击获取
返回列表