Spark与Django构建猫眼电影推荐系统实战
2026/7/24 9:58:06 网站建设 项目流程

1. 项目概述:当Spark遇上猫眼电影数据

这个项目本质上是一个融合了大数据处理与Web应用的完整数据流水线系统。我去年为本地一家影院连锁品牌实施过类似方案,核心目标是通过分析猫眼平台的电影评分、票房和用户评论数据,为影院排片和会员推荐提供数据支撑。

系统采用典型的Lambda架构设计:Spark负责离线的批量数据处理和模型训练,Django搭建实时推荐服务。这种组合既能处理海量历史数据(我们处理的原始数据量约37GB),又能保证推荐结果的低延迟响应(平均响应时间控制在120ms内)。实际运行中,每周用Spark处理新增数据,每日通过Django接口服务提供超过2万次推荐。

关键设计选择:没有选用Flask而采用Django,主要是考虑到后台管理、用户认证等企业级功能开箱即用。实测证明,Django ORM与Spark SQL的配合度超出预期。

2. 核心架构解析

2.1 数据采集层设计

猫眼数据的获取需要处理几个特殊挑战:

  • 动态加载内容需要模拟滚动操作
  • 评分数据有IP访问频率限制
  • 影片详情页URL没有明显规律

我们最终采用的方案是:

# 使用Selenium+ChromeDriver处理动态加载 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实现

推荐服务接口需要考虑的几个关键点:

  1. 冷启动问题:新用户推荐采用热度榜+类型偏好组合
  2. 实时性要求:使用Redis缓存用户最近行为
  3. 结果多样性:在推荐结果中混入10%的探索性内容

典型接口实现:

# views.py class RecommendView(APIView): def get(self, request): user_id = request.GET.get('uid') # 优先读取实时特征 recent_views = cache.lrange(f'user:{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 数据清洗中的特殊处理

猫眼数据有几个需要特别注意的清洗点:

  1. 评分标准化:将"9.5分"转换为数值9.5

    df = df.withColumn("rating", regexp_extract(col("rating_str"), "(\d+\.?\d*)", 1).cast("float"))
  2. 时间字段处理

    // 处理"上映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"))
  3. 评论情感分析: 使用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的改动:

  1. 数据库层面:

    • 使用select_relatedprefetch_related减少查询次数
    • 对电影表添加django.contrib.postgres.indexes.GinIndex
  2. 缓存策略:

    # 使用两级缓存 def get_movie_detail(movie_id): result = cache.get(f'movie:{movie_id}') if not result: result = Movie.objects.filter(...).first() cache.set(f'movie:{movie_id}', result, timeout=3600) cache.set(f'movie:{movie_id}:backup', result, timeout=86400) return result
  3. 异步任务: 使用Celery处理日志分析和推荐结果预计算:

    @app.task(bind=True) def update_recs(self, user_id): try: # 调用Spark Thrift Server conn = hive.connect(thrift_host) cursor = conn.cursor() cursor.execute(f"CALL update_user_recs({user_id})") except Exception as e: self.retry(exc=e, countdown=60)

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 60s

5.2 Django接口问题

跨域问题

CORS_ALLOWED_ORIGINS = [ "https://yourdomain.com", "http://localhost:8080" ] CORS_EXPOSE_HEADERS = ['X-Recommend-Source']

性能瓶颈排查

  1. 使用django-debug-toolbar分析SQL查询
  2. @silk_profile装饰器定位慢接口
  3. 检查Nginx和uWSGI的worker配置

6. 项目扩展方向

在实际运营中,我们发现几个有价值的扩展点:

  1. 实时推荐流

    • 接入Kafka处理用户实时行为
    • 使用Spark Streaming更新推荐结果
    val kafkaStream = KafkaUtils.createDirectStream[...] kafkaStream.foreachRDD { rdd => rdd.map(parseUserAction) .filter(_.actionType == "CLICK") .foreachPartition(updateUserProfile) }
  2. 多维度分析

    • 影院上座率预测(需接入票务数据)
    • 影片类型流行度地域分析
  3. A/B测试框架

    # 简单的分组实验实现 def get_rec_group(user_id): key = f"exp:rec:{user_id}" group = cache.get(key) if not group: group = "A" if hash(user_id) % 2 == 0 else "B" cache.set(key, group, timeout=86400*7) return group

这个项目最让我意外的发现是:周末晚间时段的用户更倾向于接受推荐(点击率比工作日高42%),我们因此调整了推荐策略的时间权重参数。大数据项目最迷人的地方就在于,数据总会给你意想不到的insight。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询