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

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)

你之前遇到的错误原因解析

  1. Window函数误用:Window函数是为每行添加开窗聚合结果,不是将同组数据合并为单行,要实现聚合合并必须用groupBy。
  2. concat_ws类型不匹配:concat_ws参数要求为字符串列,直接传入数组列会报错,需先用array_join将数组转成字符串,或用flatten合并数组后再处理。
  3. UDF类型错误:自定义UDF时,必须保证输入输出类型与DataFrame列类型严格匹配(比如处理Double类型数组时,UDF返回类型要声明正确),否则会触发类型不兼容错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 11:03:19