You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在PySpark中避免使用for循环按流派统计电影评分与数量

解决方案:用explode拆分流派数组,实现无循环统计

你当前的代码通过手动处理流派数组的每个索引来统计数据,操作繁琐且扩展性差。PySpark的explode函数可以直接将数组中的每个流派拆分为独立行,无需循环或手动索引,就能轻松实现按流派分组统计的需求。

优化后的完整代码

import pyspark
from pyspark.sql import SparkSession
from pyspark.sql.functions import split, avg, count, col, concat_ws, explode, sum

spark = SparkSession.builder.appName("APISpark").getOrCreate()

# 读取并预处理评分数据
ratings = spark.read.option("header", "true").csv("input/ml25m/ratings.csv").drop("userId", "timestamp")
# 读取电影数据并拆分流派为数组
movies = spark.read.option("header", "true").csv("input/ml-25m/movies.csv")
movies_with_genres = movies.withColumn("genre", split(movies["genres"], "\|")).drop("genres")

# 用explode将数组中的每个流派拆分为单独行
exploded_movies = movies_with_genres.select("movieId", "title", explode(col("genre")).alias("genre"))

# 关联评分数据并按流派统计
joined_data = exploded_movies.join(ratings, on="movieId", how="inner").drop("movieId")
genre_stats = joined_data.groupBy("genre")\
    .agg(
        sum("rating").alias("total_rating"),
        count("title").alias("total_count")
    )\
    .withColumn("average_rating", col("total_rating") / col("total_count"))\
    .select("genre", "average_rating", "total_count")

# 输出结果
genre_stats.select(
    concat_ws(",", col("genre"), col("average_rating"), col("total_count"))\
    .alias("genre_averagerating_Promedio_reviews")
).write.text("3_out")

关键逻辑说明

  • explode函数的作用:将多流派的数组列拆分为单流派的扁平结构。比如一部电影标注了Action|Comedy,拆分后会生成两行数据,每行对应一个流派,彻底解决了多流派的统计问题。
  • 简化分组统计:拆分后直接按genre分组,用sum计算该流派的总评分、count计算总评价数量,再通过除法得到平均分,逻辑清晰且避免了多次分组、join的冗余操作。
  • 性能提升:相比你原有的多次分组和join操作,这种方式减少了大量中间计算步骤,在大数据量场景下性能优势明显。

内容的提问来源于stack exchange,提问作者camilanahuel

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.25 10:55:12