Spark与Django构建猫眼电影推荐系统实战

Spark与Django构建猫眼电影推荐系统实战
1. 项目概述当Spark遇上猫眼电影数据这个项目本质上是一个融合了大数据处理与Web应用的完整数据流水线系统。我去年为本地一家影院连锁品牌实施过类似方案核心目标是通过分析猫眼平台的电影评分、票房和用户评论数据为影院排片和会员推荐提供数据支撑。系统采用典型的Lambda架构设计Spark负责离线的批量数据处理和模型训练Django搭建实时推荐服务。这种组合既能处理海量历史数据我们处理的原始数据量约37GB又能保证推荐结果的低延迟响应平均响应时间控制在120ms内。实际运行中每周用Spark处理新增数据每日通过Django接口服务提供超过2万次推荐。关键设计选择没有选用Flask而采用Django主要是考虑到后台管理、用户认证等企业级功能开箱即用。实测证明Django ORM与Spark SQL的配合度超出预期。2. 核心架构解析2.1 数据采集层设计猫眼数据的获取需要处理几个特殊挑战动态加载内容需要模拟滚动操作评分数据有IP访问频率限制影片详情页URL没有明显规律我们最终采用的方案是# 使用SeleniumChromeDriver处理动态加载 driver.execute_script(window.scrollTo(0, document.body.scrollHeight);) time.sleep(random.uniform(1.5, 3)) # 随机延时规避反爬 # 分布式爬虫架构 scrapy_redis 阿里云函数计算实现IP自动切换2.2 Spark数据处理流水线数据处理阶段最耗时的操作是用户-电影评分矩阵的构建。这里采用了Spark的优化技巧// 使用ALS算法时的参数优化 val als new ALS() .setRank(50) // 隐语义维度 .setMaxIter(15) // 迭代次数 .setRegParam(0.01) // 正则化参数 .setUserCol(userId) .setItemCol(movieId) .setRatingCol(rating) // 特别重要的缓存策略 val ratings spark.read.parquet(...) .repartition(200) // 根据集群核数调整 .persist(StorageLevel.MEMORY_AND_DISK_SER)2.3 Django推荐API实现推荐服务接口需要考虑的几个关键点冷启动问题新用户推荐采用热度榜类型偏好组合实时性要求使用Redis缓存用户最近行为结果多样性在推荐结果中混入10%的探索性内容典型接口实现# views.py class RecommendView(APIView): def get(self, request): user_id request.GET.get(uid) # 优先读取实时特征 recent_views cache.lrange(fuser:{user_id}:recent, 0, 4) # 混合推荐逻辑 if len(recent_views) 2: recs spark_client.get_cf_recs(user_id) # 协同过滤结果 else: recs get_trending_movies() # 热门电影 # 添加多样性 if random.random() 0.1: recs[-1] get_random_movie() return Response(recs)3. 关键技术实现细节3.1 数据清洗中的特殊处理猫眼数据有几个需要特别注意的清洗点评分标准化将9.5分转换为数值9.5df df.withColumn(rating, regexp_extract(col(rating_str), (\d\.?\d*), 1).cast(float))时间字段处理// 处理上映3天这类相对时间 val releaseDate when(col(date_str).contains(天), date_sub(current_date(), regexp_extract(col(date_str), (\d), 1).cast(int))) .otherwise(to_date(col(date_str), yyyy-MM-dd))评论情感分析 使用HanLP自定义电影领域词典准确率提升23%from pyhanlp import * analyzer PerceptronLexicalAnalyzer() analyzer.enableCustomDictionaryForcing(True)3.2 推荐算法优化经过AB测试最终采用的混合推荐策略算法类型使用场景准确率覆盖率ALS协同过滤老用户推荐0.720.65内容相似度新电影推荐0.680.82热度加权冷启动阶段0.610.95关键优化点为ALS添加时间衰减因子weight 1 / (1 log(1 days_ago))内容特征使用BERT向量而非TF-IDF实时点击行为影响权重设为离线数据的1.8倍4. 部署与性能调优4.1 Spark集群配置在8节点集群上的最优配置每节点16核64GB# spark-defaults.conf关键配置 spark.executor.memory 48G spark.executor.cores 12 spark.driver.memory 8G spark.sql.shuffle.partitions 600 spark.default.parallelism 400 spark.serializer org.apache.spark.serializer.KryoSerializer重要教训spark.sql.shuffle.partitions设置过小会导致OOM过大则降低效率。建议设为集群总核数的2-3倍。4.2 Django性能优化几个显著提升QPS的改动数据库层面使用select_related和prefetch_related减少查询次数对电影表添加django.contrib.postgres.indexes.GinIndex缓存策略# 使用两级缓存 def get_movie_detail(movie_id): result cache.get(fmovie:{movie_id}) if not result: result Movie.objects.filter(...).first() cache.set(fmovie:{movie_id}, result, timeout3600) cache.set(fmovie:{movie_id}:backup, result, timeout86400) return result异步任务 使用Celery处理日志分析和推荐结果预计算app.task(bindTrue) def update_recs(self, user_id): try: # 调用Spark Thrift Server conn hive.connect(thrift_host) cursor conn.cursor() cursor.execute(fCALL update_user_recs({user_id})) except Exception as e: self.retry(exce, countdown60)5. 典型问题排查实录5.1 Spark常见报错处理问题1Container killed by YARN for exceeding memory limits解决方案检查executor内存分配是否合理添加spark.executor.memoryOverhead建议设为executor内存的10-15%对大数据集使用persist(StorageLevel.MEMORY_AND_DISK_SER)问题2java.net.SocketTimeoutException: Read timed out处理方法# 增加超时阈值 spark.network.timeout 600s spark.executor.heartbeatInterval 60s5.2 Django接口问题跨域问题CORS_ALLOWED_ORIGINS [ https://yourdomain.com, http://localhost:8080 ] CORS_EXPOSE_HEADERS [X-Recommend-Source]性能瓶颈排查使用django-debug-toolbar分析SQL查询用silk_profile装饰器定位慢接口检查Nginx和uWSGI的worker配置6. 项目扩展方向在实际运营中我们发现几个有价值的扩展点实时推荐流接入Kafka处理用户实时行为使用Spark Streaming更新推荐结果val kafkaStream KafkaUtils.createDirectStream[...] kafkaStream.foreachRDD { rdd rdd.map(parseUserAction) .filter(_.actionType CLICK) .foreachPartition(updateUserProfile) }多维度分析影院上座率预测需接入票务数据影片类型流行度地域分析A/B测试框架# 简单的分组实验实现 def get_rec_group(user_id): key fexp:rec:{user_id} group cache.get(key) if not group: group A if hash(user_id) % 2 0 else B cache.set(key, group, timeout86400*7) return group这个项目最让我意外的发现是周末晚间时段的用户更倾向于接受推荐点击率比工作日高42%我们因此调整了推荐策略的时间权重参数。大数据项目最迷人的地方就在于数据总会给你意想不到的insight。