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实现
推荐服务接口需要考虑的几个关键点:
- 冷启动问题:新用户推荐采用热度榜+类型偏好组合
- 实时性要求:使用Redis缓存用户最近行为
- 结果多样性:在推荐结果中混入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 数据清洗中的特殊处理
猫眼数据有几个需要特别注意的清洗点:
评分标准化:将"9.5分"转换为数值9.5
df = 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.72 | 0.65 |
| 内容相似度 | 新电影推荐 | 0.68 | 0.82 |
| 热度加权 | 冷启动阶段 | 0.61 | 0.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(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异步任务: 使用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常见报错处理
问题1:Container killed by YARN for exceeding memory limits
解决方案:
- 检查executor内存分配是否合理
- 添加
spark.executor.memoryOverhead(建议设为executor内存的10-15%) - 对大数据集使用
persist(StorageLevel.MEMORY_AND_DISK_SER)
问题2:java.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 = 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。