如何使用Spark DataFrame创建以十年为单位的窗口?
如何使用Spark DataFrame创建以十年为单位的窗口?
我来帮你一步步搞定这个需求!要实现按十年分组的窗口,核心思路是先把电影的年份转换成对应的十年分组标识,再基于这个标识创建窗口。下面结合你给的示例数据,用Python版Spark来演示具体操作:
1. 定义Schema并加载示例数据
首先,我们先明确你提供的数据集字段:年份;ID;电影名;类型;演员1;演员2;导演;评分;未知字段;图片名,先定义对应的Schema,再加载数据:
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, IntegerType, StringType # 初始化SparkSession spark = SparkSession.builder.appName("DecadeWindowExample").getOrCreate() # 定义Schema schema = StructType([ StructField("year", IntegerType(), nullable=False), StructField("id", IntegerType(), nullable=False), StructField("title", StringType(), nullable=False), StructField("genre", StringType(), nullable=False), StructField("actor1", StringType(), nullable=False), StructField("actor2", StringType(), nullable=False), StructField("director", StringType(), nullable=False), StructField("rating", IntegerType(), nullable=False), StructField("unknown_col", StringType(), nullable=True), StructField("image", StringType(), nullable=True) ]) # 示例数据(模拟你提供的内容) sample_data = [ "1990;111;Tie Me Up! Tie Me Down!;Comedy;Banderas, Antonio;Abril, Victoria;Almodóvar, Pedro;68;No;NicholasCage.png", "1991;113;High Heels;Comedy;Bosé, Miguel;Abril, Victoria;Almodóvar, Pedro;68;No;NicholasCage.png", "1983;104;Dead Zone, The;Horror;Walken, Christopher;Adams, Brooke;Cronenberg, David;79;No;NicholasCage.png", "1979;122;Cuba;Action;Connery, Sean;Adams, Brooke;Lester, Richard;6;No;seanConnery.png", "1978;94;Days of Heaven;Drama;Gere, Richard;Adams, Brooke;Malick, Terrence;14;No;NicholasCage.png" ] # 加载数据 df = spark.createDataFrame([row.split(";") for row in sample_data], schema=schema) df.show(truncate=False)
2. 添加十年分组列
我们用整数除法把年份转换成十年的起始年份(比如1990、1991都变成1990,代表90年代):
from pyspark.sql.functions import col # 计算十年分组:(年份 // 10) * 10 得到十年起始年 df_with_decade = df.withColumn("decade", (col("year") // 10) * 10) df_with_decade.select("year", "decade", "title").show()
执行后你会看到:
+----+------+---------------------------+ |year|decade|title | +----+------+---------------------------+ |1990|1990 |Tie Me Up! Tie Me Down! | |1991|1990 |High Heels | |1983|1980 |Dead Zone, The | |1979|1970 |Cuba | |1978|1970 |Days of Heaven | +----+------+---------------------------+
3. 创建以十年为单位的窗口
现在我们可以基于decade字段创建窗口了,如果你需要在十年内按年份排序,可以加上orderBy:
from pyspark.sql.window import Window # 创建窗口:按十年分组,组内按年份升序排列 decade_window = Window.partitionBy("decade").orderBy(col("year").asc())
4. 窗口函数的使用示例
有了窗口后,我们可以做各种操作,比如:
- 给每个十年内的电影按年份排序加行号
- 计算每个十年的平均评分
- 统计每个十年的电影数量
示例1:十年内电影按年份排序的行号
from pyspark.sql.functions import row_number df_with_row_num = df_with_decade.withColumn("row_num_in_decade", row_number().over(decade_window)) df_with_row_num.select("decade", "year", "title", "row_num_in_decade").show()
输出:
+------+----+---------------------------+------------------+ |decade|year|title |row_num_in_decade | +------+----+---------------------------+------------------+ |1970 |1978|Days of Heaven |1 | |1970 |1979|Cuba |2 | |1980 |1983|Dead Zone, The |1 | |1990 |1990|Tie Me Up! Tie Me Down! |1 | |1990 |1991|High Heels |2 | +------+----+---------------------------+------------------+
示例2:计算每个十年的平均评分(聚合窗口)
如果需要计算每个十年的整体统计值,可以用聚合窗口:
from pyspark.sql.functions import avg # 聚合窗口:只按十年分组,不排序 decade_agg_window = Window.partitionBy("decade") df_with_avg_rating = df_with_decade.withColumn("avg_rating_in_decade", avg(col("rating")).over(decade_agg_window)) df_with_avg_rating.select("decade", "title", "rating", "avg_rating_in_decade").show()
输出:
+------+---------------------------+------+---------------------+ |decade|title |rating|avg_rating_in_decade | +------+---------------------------+------+---------------------+ |1970 |Days of Heaven |14 |10.0 | |1970 |Cuba |6 |10.0 | |1980 |Dead Zone, The |79 |79.0 | |1990 |Tie Me Up! Tie Me Down! |68 |68.0 | |1990 |High Heels |68 |68.0 | +------+---------------------------+------+---------------------+
这样就完成了以十年为单位的窗口创建和使用啦!如果是用Scala的话,逻辑是完全一样的,只是语法稍有不同~
内容的提问来源于stack exchange,提问作者prady
相关产品推荐
相关产品推荐

