PySpark按周/月聚合列表类型列的实现问题求助
PySpark/Databricks按周/月聚合拼接列表列实现方案
需求明确
将DataFrame中brands(字符串数组)、weight(数值数组)列,按指定周区间(如2023-04-01 至 2023-04-07)或月区间聚合,把同周期内的所有列表元素合并为单行的数组或字符串。
步骤1:构造示例测试数据
先创建模拟场景的DataFrame,方便验证逻辑:
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.types import StructType, StructField, DateType, ArrayType, StringType, DoubleType spark = SparkSession.builder.appName("AggregateListColumns").getOrCreate() data = [ ("2023-04-02", ["A", "B"], [1.2, 3.4]), ("2023-04-05", ["C"], [2.5]), ("2023-04-09", ["A", "D"], [0.8, 4.1]), ("2023-04-15", ["B", "C"], [1.9, 2.2]), ("2023-05-01", ["A"], [3.0]) ] schema = StructType([ StructField("date", DateType(), True), StructField("brands", ArrayType(StringType()), True), StructField("weight", ArrayType(DoubleType()), True) ]) df = spark.createDataFrame(data, schema) df.show(truncate=False)
步骤2:生成周/月分组区间
把日期转换为需求中的周、月区间字符串,作为聚合的分组键:
# 生成周区间(每周从周一到周日,格式:YYYY-MM-DD 至 YYYY-MM-DD) df_with_periods = df.withColumn( "week_period", F.concat( F.date_format(F.date_sub(F.next_day(F.col("date"), "Sunday"), 6), "yyyy-MM-dd"), F.lit(" 至 "), F.date_format(F.next_day(F.col("date"), "Sunday"), "yyyy-MM-dd") ) ).withColumn( "month_period", F.concat( F.date_format(F.date_trunc("month", F.col("date")), "yyyy-MM-dd"), F.lit(" 至 "), F.date_format(F.last_day(F.col("date")), "yyyy-MM-dd") ) ) df_with_periods.show(truncate=False)
步骤3:按周/月聚合拼接列表
场景1:合并为大数组后转字符串
将同周期内的brands、weight分别合并成大数组,再拼接成逗号分隔的字符串:
# 按周聚合 weekly_agg = df_with_periods.groupBy("week_period").agg( # 合并所有行的brands数组为一个大数组 F.flatten(F.collect_list("brands")).alias("all_brands"), # 合并所有行的weight数组为一个大数组 F.flatten(F.collect_list("weight")).alias("all_weights") ).withColumn( "brands_str", F.array_join(F.col("all_brands"), ", ") # 数组转字符串 ).withColumn( "weights_str", F.array_join(F.col("all_weights").cast(ArrayType(StringType())), ", ") # 数值数组先转字符串数组再拼接 ) weekly_agg.show(truncate=False) # 按月聚合(逻辑与周一致) monthly_agg = df_with_periods.groupBy("month_period").agg( F.flatten(F.collect_list("brands")).alias("all_brands"), F.flatten(F.collect_list("weight")).alias("all_weights") ).withColumn( "brands_str", F.array_join(F.col("all_brands"), ", ") ).withColumn( "weights_str", F.array_join(F.col("all_weights").cast(ArrayType(StringType())), ", ") ) monthly_agg.show(truncate=False)
场景2:保持brands与weight的对应关系
如果需要每个品牌对应其权重(如A:1.2, B:3.4),先将每行的brands和weight配对,再聚合:
# 先将每行的brands和weight配对成结构体数组 df_with_pairs = df_with_periods.withColumn( "brand_weight_pair", F.arrays_zip(F.col("brands"), F.col("weight")) ) # 按周聚合配对后的数组,再转成字符串 weekly_agg_pairs = df_with_pairs.groupBy("week_period").agg( F.flatten(F.collect_list("brand_weight_pair")).alias("all_brand_weight_pairs") ).withColumn( "brand_weight_str", F.array_join( F.transform( F.col("all_brand_weight_pairs"), lambda x: F.concat(x.brands, F.lit(":"), x.weight.cast(StringType())) ), ", " ) ) weekly_agg_pairs.show(truncate=False)
你之前遇到的错误原因解析
- Window函数误用:Window函数是为每行添加开窗聚合结果,不是将同组数据合并为单行,要实现聚合合并必须用
groupBy。 - concat_ws类型不匹配:
concat_ws参数要求为字符串列,直接传入数组列会报错,需先用array_join将数组转成字符串,或用flatten合并数组后再处理。 - UDF类型错误:自定义UDF时,必须保证输入输出类型与DataFrame列类型严格匹配(比如处理Double类型数组时,UDF返回类型要声明正确),否则会触发类型不兼容错误。
内容的提问来源于stack exchange,提问作者user1717931
相关产品推荐
相关产品推荐

