PySpark中向Struct类型字段添加已有字段并分组聚合问题
解决PySpark DataFrame字段嵌入数组Struct并分组聚合的问题
原始数据与Schema
输入DataFrame:
sku category price state infos_gerais 33344 mmmma 3.00 SP [{5, 5656655, 5845454}] 33344 mmmma 3.00 MG [{5, 6565767, 5854545}] 33344 mmmma 3.00 RS [{5, 8788787, 4564646}]
数据Schema:
|-- sku: string (nullable = true) |-- category: string (nullable = true) |-- price: double (nullable = true) |-- state: string (nullable = true) |-- infos_gerais: array (nullable = true) | |-- element: struct (containsNull = false) | | |-- service_type_id: integer (nullable = true) | | |-- cep_ini: integer (nullable = true) | | |-- cep_fim: integer (nullable = true)
需求说明
将顶层state字段添加到infos_gerais数组的每个Struct元素中,再按sku、category、price分组聚合,最终合并所有Struct元素到一个数组。
原代码问题分析
你提供的代码存在两个核心问题:
- 错误引用
infos_gerais.state:state是顶层字段,并非infos_gerais数组内的Struct字段 collect_list是聚合函数,不能直接在withColumn中使用,必须配合groupBy操作
正确实现代码
from pyspark.sql import functions as sf # 方案1:先处理数组添加字段,再展开聚合 df_processed = df.withColumn( "infos_gerais", sf.transform( "infos_gerais", lambda x: sf.struct(x["service_type_id"], x["cep_ini"], x["cep_fim"], sf.col("state").alias("state")) ) ) df_result = df_processed.select( "sku", "category", "price", sf.explode("infos_gerais").alias("infos") ).groupBy("sku", "category", "price").agg( sf.collect_list("infos").alias("infos_gerais") ) # 方案2:直接在聚合逻辑中完成字段添加与合并 df_result = df.groupBy("sku", "category", "price").agg( sf.flatten( sf.collect_list( sf.transform( "infos_gerais", lambda x: sf.struct(x["service_type_id"], x["cep_ini"], x["cep_fim"], sf.col("state").alias("state")) ) ) ).alias("infos_gerais") ) # 查看结果 df_result.show(truncate=False)
期望输出
sku category price infos_gerais 33344 mmmma 3.00 [{5, 5656655, 5845454, SP}, {5, 6565767, 5854545, MG},{5, 8788787, 4564646, RS}]
内容的提问来源于stack exchange,提问作者Vivian
相关产品推荐
相关产品推荐

