拓冰建站拓冰建站
首页 / 资讯中心 / 正文

Python+Spark+Hadoop构建生产级电影推荐系统

简介本资源是一套基于Python、Spark与Hadoop构建的分布式电影推荐系统毕业设计源码面向大数据与人工智能方向的本科生及初学者解决用户画像建模与个性化推荐落地实践问题。压缩包共801个文件含60个核心Python脚本含数据预处理、Spark任务提交、推荐算法实现、340个前端JS/CSS文件支撑可视化管理界面、151个CSS样式资源及21个HTML页面另有SQL建表语句、日志与说明文档等整体16.21MB结构完整覆盖后端计算、前端展示与数据流转全流程。已有890人学习下载提供可直接运行的端到端工程包含HDFS数据存储配置、Spark MLlib协同过滤实现、用户标签体系构建逻辑、以及基于行为日志生成画像的完整链路代码适合用于课程设计复现、毕设参考或分布式推荐系统入门实战。1. 用 Python Spark Hadoop 搭建可复现的电影推荐系统不是跑通 demo而是构建带用户画像闭环的生产级数据链路你下载了一个名为“PythonSparkHadoop大数据基于用户画像电影推荐系统毕业源码 - 副本.zip”的压缩包解压后发现目录结构混乱、配置文件缺失、README 里只有一句“运行 main.py 即可”结果 pip install 后报 ModuleNotFoundError: No module named pyspark再查 spark-submit 报错 class not found最后在本地伪分布式 Hadoop 上连 hdfs dfs -ls 都超时——这不是代码问题而是整条数据栈的衔接断层。这个标题指向的不是一个“能跑的毕设”而是一套必须同时满足三重约束的真实场景Python 承担特征工程与服务封装、Spark 负责高吞吐画像计算与协同过滤训练、Hadoop 提供稳定可靠的原始日志存储与中间表管理。它适合正在从单机 Pandas 过渡到分布式计算的中级工程师也适合需要快速验证用户分群推荐效果的数据产品同学。关键不在“有没有推荐”而在“画像标签能否回溯、推荐结果能否 AB 测试、模型更新能否触发 HDFS 数据自动归档”——这些能力全藏在 Spark 读写 HDFS 的路径约定、用户行为日志的分区设计、以及 Python 推荐服务调用 Spark MLlib 模型的序列化方式里。2. 用户画像构建从原始日志到可查询标签体系用 Spark SQL 实现分层加工用户画像是整个系统的锚点不是简单统计“用户看了什么”而是建立可扩展、可验证、可下钻的标签体系。常见误区是直接用 Python pandas 处理 CSV 日志但当用户行为日志日增 500 万条约 2GB单机内存和 IO 成为瓶颈。此时必须依赖 Spark 在 HDFS 上的分布式计算能力且标签加工需遵循分层设计原则ODS 层存原始日志、DWD 层做清洗与原子事件标准化、DWS 层聚合用户维度宽表、ADS 层输出业务可用标签。2.1 原始日志接入 HDFS 并按天分区确保 Spark 可高效扫描假设原始日志为 JSON 格式每行一条用户行为记录字段包括user_id,movie_id,rating,timestamp,device_type。首先需将日志文件上传至 HDFS 指定路径并严格按日期分区# 创建 HDFS 目录并设置权限Hadoop 用户执行 hdfs dfs -mkdir -p /data/movie_logs/dt2024-06-01 hdfs dfs -put ./raw_logs_20240601.json /data/movie_logs/dt2024-06-01/ hdfs dfs -chmod -R 755 /data/movie_logs注意分区路径必须为dtYYYY-MM-DD格式这是 Spark SQL 自动识别分区的关键。不要用/data/movie_logs/2024/06/01/这类嵌套路径否则后续MSCK REPAIR TABLE无法自动加载分区。2.2 用 Spark SQL 构建 DWD 层原子事件表统一时间戳与设备归一在 PySpark 中启动 SparkSession连接到已配置好 core-site.xml 和 hdfs-site.xml 的集群本地伪分布式或 YARN 模式均可from pyspark.sql import SparkSession from pyspark.sql.functions import from_unixtime, col, when, lit from pyspark.sql.types import StructType, StructField, StringType, IntegerType, TimestampType spark SparkSession.builder \ .appName(movie_dwd_builder) \ .config(spark.sql.adaptive.enabled, true) \ .config(spark.sql.hive.metastore.uris, thrift://localhost:9083) \ .enableHiveSupport() \ .getOrCreate() # 定义 schema 显式解析 JSON避免推断错误 schema StructType([ StructField(user_id, StringType(), True), StructField(movie_id, StringType(), True), StructField(rating, IntegerType(), True), StructField(timestamp, StringType(), True), # 原始为字符串时间戳 StructField(device_type, StringType(), True) ]) # 读取分区数据自动识别 dt2024-06-01 dwd_df spark.read \ .schema(schema) \ .json(hdfs://localhost:9000/data/movie_logs) \ .withColumn(event_time, from_unixtime(col(timestamp).cast(long))) \ .withColumn(device_category, when(col(device_type).isin([ios, android]), mobile) .when(col(device_type).isin([pc, mac]), desktop) .otherwise(other)) \ .filter(col(user_id).isNotNull() col(movie_id).isNotNull()) # 写入 Hive 表需提前建库 create database movie_dw dwd_df.write \ .mode(overwrite) \ .partitionBy(dt) \ .format(parquet) \ .saveAsTable(movie_dw.dwd_user_behavior)这段代码完成三件事强制 schema 解析避免空值崩溃、统一时间格式便于后续窗口计算、归一化设备类型降低下游判断复杂度。saveAsTable将数据持久化到 Hive Metastore后续所有 SQL 查询可直接引用movie_dw.dwd_user_behavior表名无需关心底层 HDFS 路径。2.3 DWS 层构建用户画像宽表行为频次、偏好强度、活跃度分层画像宽表是推荐模型的直接输入源。我们不直接用原始行为做推荐而是提取稳定、低噪声的聚合特征-- 在 beeline 或 Spark SQL CLI 中执行 CREATE TABLE movie_dw.dws_user_profile AS SELECT user_id, COUNT(*) AS total_views, COUNT(DISTINCT movie_id) AS distinct_movies, AVG(rating) AS avg_rating, -- 最近7天活跃度以登录或播放行为为准 SUM(CASE WHEN event_time date_sub(current_date(), 7) THEN 1 ELSE 0 END) AS recent_active_days, -- 类型偏好统计各类型电影观看占比需关联 movie_info 表 COLLECT_LIST( STRUCT( genre, COUNT(*) AS genre_cnt ) ) AS genre_preference FROM movie_dw.dwd_user_behavior a JOIN movie_dw.dim_movie_info b ON a.movie_id b.movie_id GROUP BY user_id;提示COLLECT_LIST(STRUCT(...))是 Spark SQL 中生成嵌套结构的常用手法比用 UDF 更高效。genre_preference字段类型为arraystructgenre:string,genre_cnt:bigintPython 侧可通过row.genre_preference[0].genre直接访问无需 JSON 解析。3. 推荐模型训练与部署用 Spark MLlib 训练 ALS 模型导出为 Python 可加载格式电影推荐本质是矩阵分解问题ALSAlternating Least Squares算法在 Spark MLlib 中成熟稳定、支持隐式反馈如点击、播放时长、且天然适配分布式训练。关键不在“调用 fit()”而在如何让训练好的模型脱离 Spark 环境被 Python Web 服务实时调用。3.1 使用隐式反馈训练 ALS 模型规避评分稀疏性陷阱显式评分1~5 分在真实场景中占比极低更多是播放完成率、停留时长等隐式信号。我们将watch_duration_sec作为置信度权重from pyspark.ml.recommendation import ALS from pyspark.sql.functions import col, when, log # 从 DWD 层读取行为数据构造隐式反馈 implicit_df spark.table(movie_dw.dwd_user_behavior) \ .withColumn(confidence, when(col(rating) 0, 10.0 log(col(rating))) # 显式评分加权 .otherwise(1.0 col(watch_duration_sec) / 100.0)) # 隐式时长加权 als ALS( maxIter10, regParam0.01, rank50, # 潜在因子数50~100 为常用范围 userColuser_id, itemColmovie_id, ratingColconfidence, coldStartStrategydrop, # 对新用户/新电影不预测避免 NaN nonnegativeTrue ) model als.fit(implicit_df) # 保存模型到 HDFS路径需可被 Python 进程访问 model.write().overwrite().save(hdfs://localhost:9000/model/als_model_v20240601)coldStartStrategydrop是关键参数生产环境绝不允许模型对未见过的用户返回随机推荐必须明确拒绝。nonnegativeTrue强制因子向量非负提升可解释性。3.2 导出用户/物品因子矩阵为 Parquet供 Python 直接读取Spark MLlib 模型本身无法被 Python pickle 加载但其核心——用户因子矩阵userFactors和物品因子矩阵itemFactors——可导出为标准 Parquet# 导出用户因子 model.userFactors.write.mode(overwrite).parquet(hdfs://localhost:9000/model/user_factors_v20240601) # 导出物品因子 model.itemFactors.write.mode(overwrite).parquet(hdfs://localhost:9000/model/item_factors_v20240601)在 Python 服务端使用pandas.read_parquet()或pyarrow.parquet.read_table()即可加载import pandas as pd import pyarrow.parquet as pq # 从 HDFS 下载或挂载 NFS 后读取 user_factors_df pd.read_parquet(/mnt/hdfs/model/user_factors_v20240601) item_factors_df pd.read_parquet(/mnt/hdfs/model/item_factors_v20240601) # 转为 numpy 矩阵用于余弦相似度计算 import numpy as np U user_factors_df.set_index(id)[features].apply(lambda x: np.array(x)).tolist() I item_factors_df.set_index(id)[features].apply(lambda x: np.array(x)).tolist() U_matrix np.vstack(U) I_matrix np.vstack(I)注意features列是 Spark 生成的Vector类型Parquet 中序列化为二进制数组。pd.read_parquet()会自动还原为 Python list再转np.array即可参与计算。此方式绕过 SparkContext 依赖使推荐服务彻底轻量化。3.3 实现基于因子内积的实时推荐支持 Top-N 与冷启动兜底最终推荐逻辑在 Python 中实现不调用任何 Spark 组件def get_recommendations(user_id: str, top_k: int 10) - list: # 1. 查找用户因子向量 if user_id not in user_factors_df[id].values: # 冷启动返回热门电影从 HDFS 读取预计算的 popularity.csv return get_popular_movies(top_k) user_vec user_factors_df[user_factors_df[id] user_id][features].iloc[0] user_vec np.array(user_vec) # 2. 计算与所有电影的内积即预测得分 scores I_matrix.dot(user_vec) # (n_items,) 向量 # 3. 排序并返回 movie_id 列表 top_indices np.argsort(scores)[-top_k:][::-1] return item_factors_df.iloc[top_indices][id].tolist() # 示例调用 print(get_recommendations(user_12345, top_k5)) # 输出: [movie_789, movie_456, movie_123, ...]该函数耗时 50ms在 10 万电影规模下完全满足 API 实时响应要求。get_popular_movies()可从 HDFS 读取每日更新的热门榜 CSV实现无模型兜底。4. Hadoop 与 Spark 协同配置解决本地开发与集群提交的路径一致性问题本地调试成功但spark-submit到 YARN 集群时提示java.io.IOException: Failed on local exception: java.io.IOException: Rejected socket connection——这通常不是代码错误而是 Hadoop 和 Spark 的资源配置未对齐。核心矛盾在于Python 侧读写 HDFS 的路径写法、Spark 作业提交时的 master 配置、以及 HDFS NameNode 地址在不同环境下的映射关系。4.1 统一 HDFS 访问协议用hdfs://替代file://并配置 core-site.xml无论本地伪分布式还是集群模式所有路径必须以hdfs://开头。在core-site.xml中明确指定 NameNode 地址!-- $HADOOP_HOME/etc/hadoop/core-site.xml -- configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value !-- 伪分布式用 localhost -- !-- 生产集群则改为 hdfs://mycluster -- /property /configurationPySpark 会自动读取该配置。若在代码中硬编码hdfs://localhost:9000当部署到集群时需修改代码——这是反模式。正确做法是只写相对路径由配置驱动# ✅ 正确依赖 core-site.xml spark.read.parquet(/model/user_factors_v20240601) # ❌ 错误硬编码 host:port spark.read.parquet(hdfs://localhost:9000/model/user_factors_v20240601)4.2 Spark 提交时指定 deploy-mode 与 driver-memory避免 OOM本地开发常忽略资源参数导致集群提交失败。关键参数表参数本地开发建议值YARN 集群建议值说明--masterlocal[*]yarn本地用所有 CPU 核集群走 YARN 调度--deploy-modeclientclusterclient 模式 driver 在本地cluster 模式 driver 在集群--driver-memory2g4gdriver 需加载模型矩阵内存不足会 GC 失败--executor-memory2g8g每个 executor 处理分区数据推荐 4~8g--conf spark.sql.adaptive.enabledtrue必开必开自适应查询优化显著提升 Join 性能提交命令示例spark-submit \ --master yarn \ --deploy-mode cluster \ --driver-memory 4g \ --executor-memory 8g \ --conf spark.sql.adaptive.enabledtrue \ --conf spark.sql.adaptive.coalescePartitions.enabledtrue \ movie_recommender.py提示spark.sql.adaptive.coalescePartitions.enabledtrue可自动合并小分区避免大量 task 启动开销在处理用户画像宽表这类宽列数据时尤为有效。4.3 验证 HDFS 与 Spark 连通性三步快速定位网络与权限问题当hdfs dfs -ls成功但 Spark 读取失败按顺序排查确认 Spark 能解析 HDFS URI在 PySpark Shell 中执行spark.sparkContext._jsc.hadoopConfiguration().get(fs.defaultFS) # 应返回 hdfs://localhost:9000检查 HDFS 文件权限与 ownerhdfs dfs -ls /model/user_factors_v20240601 # 确保权限为 -rwxr-xr-xowner 为 spark 或 hadoop 用户 hdfs dfs -chown spark:spark /model/user_factors_v20240601验证 Spark 能列出 HDFS 目录内容# 在 spark-submit 的 driver 中执行 sc._jvm.org.apache.hadoop.fs.FileSystem.get(spark.sparkContext._jsc.hadoopConfiguration()) \ .listStatus(spark.sparkContext._jvm.org.apache.hadoop.fs.Path(/model)) # 若抛异常则是 Kerberos 或防火墙问题5. 用户画像标签验证与推荐效果评估用 Python 脚本自动化校验数据质量模型上线后最危险的不是推荐不准而是画像标签持续漂移——比如“最近7天活跃用户”统计口径被上游修改但下游推荐未感知。必须建立自动化校验机制不依赖人工抽查。5.1 编写画像标签一致性校验脚本每日定时运行校验核心逻辑对比 DWS 层宽表与 DWD 层原子表的聚合结果是否匹配。例如验证total_views是否等于 DWD 层对应用户的行数from pyspark.sql import SparkSession from pyspark.sql.functions import count, col spark SparkSession.builder.appName(profile_validation).getOrCreate() # 读取 DWS 宽表 dws_df spark.table(movie_dw.dws_user_profile) # 读取 DWD 原子表按 user_id 聚合 dwd_agg_df spark.table(movie_dw.dwd_user_behavior) \ .groupBy(user_id) \ .agg(count(*).alias(dwd_total_views)) # 关联校验 validation_df dws_df.join(dwd_agg_df, user_id, left) \ .withColumn(match, col(total_views) col(dwd_total_views)) \ .filter(col(match) False) # 输出不一致记录仅取前10条用于告警 mismatch_count validation_df.count() if mismatch_count 0: print(f❌ 发现 {mismatch_count} 条用户画像不一致记录) validation_df.select(user_id, total_views, dwd_total_views).show(10, truncateFalse) else: print(✅ 用户画像聚合逻辑校验通过)将此脚本加入 crontab每日凌晨 2 点执行输出结果写入日志并触发企业微信告警。5.2 推荐结果多样性评估用 Gini 不纯度量化推荐列表分布Top-N 推荐易陷入“马太效应”热门电影扎堆出现。用 Gini 系数衡量推荐列表的品类分布均衡性import numpy as np from collections import Counter def calculate_gini_diversity(recommended_movie_ids: list, movie_genre_map: dict) - float: movie_genre_map: {movie_id: genre_name} 返回 Gini 系数越接近 0 越均衡越接近 1 越集中 genres [movie_genre_map.get(mid, unknown) for mid in recommended_movie_ids] genre_counts Counter(genres) n len(recommended_movie_ids) if n 0: return 0.0 # Gini 1 - Σ(pi)^2 gini 1.0 for count in genre_counts.values(): p_i count / n gini - p_i ** 2 return round(gini, 3) # 示例对 100 个用户推荐结果计算平均多样性 sample_users [user_001, user_002, ...] all_ginis [] for uid in sample_users: recs get_recommendations(uid, top_k10) gini calculate_gini_diversity(recs, genre_map) all_ginis.append(gini) print(f平均推荐多样性 Gini 系数: {np.mean(all_ginis):.3f}) # 健康阈值 0.6 为良好 0.4 需优化如引入多样性重排序Gini 系数直观反映推荐结果是否过度集中于少数类型。当该值持续低于 0.4说明 ALS 模型陷入局部最优应引入基于图神经网络的多样性增强模块或在排序阶段加入 MMFMaximal Marginal Relevance重打分。5.3 构建最小可行监控看板用 Flask Plotly 快速展示画像与推荐指标无需 Grafana一个 50 行 Flask 应用即可监控核心指标from flask import Flask, render_template import plotly.express as px import pandas as pd app Flask(__name__) app.route(/dashboard) def dashboard(): # 从 HDFS 读取最新画像统计每日更新 CSV profile_stats pd.read_csv(/mnt/hdfs/report/profile_daily_stats.csv) # 绘制活跃用户趋势 fig1 px.line(profile_stats, xdate, yactive_users, title日活跃用户数) # 绘制推荐多样性趋势 fig2 px.line(profile_stats, xdate, ygini_diversity, title推荐多样性 Gini 系数) return render_template(dashboard.html, plot1fig1.to_html(full_htmlFalse), plot2fig2.to_html(full_htmlFalse)) if __name__ __main__: app.run(host0.0.0.0, port5000)配合 Nginx 反向代理即可对外提供轻量级监控页面。关键在于profile_daily_stats.csv由前述校验脚本每日生成并写入 HDFS形成闭环。提示所有监控数据必须来自 HDFS而非本地文件。这样当服务迁移到 Kubernetes 时只需挂载 HDFS PVC无需修改代码。本文还有配套的精品资源点击获取
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门