Spark DataFrame数组行字段提取为顶级列的技术求助
解决方案
步骤1:展开结构体字段
首先提取data结构体中需要的字段,忽略无关的kafka_offset:
from pyspark.sql import functions as F # 提取data结构体中的目标字段 df = df.select( "data.date", "data.id", "data.specifications", "data.begin", "data.end" )
步骤2:将规格数组转换为Map
利用map_from_entries函数,把specifications数组中的每个元素转换成(specname, specvalue)的键值对Map:
# 转换数组为Map,键为specname,值为specvalue df = df.withColumn( "spec_map", F.map_from_entries( F.transform( "specifications", lambda item: F.struct(item.specname.alias("key"), item.specvalue.alias("value")) ) ) )
步骤3:动态提取Map中的列
先获取所有唯一的规格名称,再从Map中提取对应值作为列,并按需求重命名为spec_小写名称格式:
# 获取所有唯一的specname spec_columns = df.select(F.explode("spec_map").alias("key", "val")) \ .select("key").distinct().rdd.flatMap(lambda x: x).collect() # 逐个提取列并命名 for col_name in spec_columns: df = df.withColumn(f"spec_{col_name.lower()}", F.col("spec_map")[col_name])
步骤4:清理临时列并转换时间字段(可选)
移除不再需要的临时列,如果begin和end是时间戳类型,可转换为日期字符串:
# 移除临时数组和Map列 df = df.drop("specifications", "spec_map") # 若begin/end为时间戳,转换为日期格式(根据实际需求调整) df = df.withColumn("begin", F.from_unixtime("begin")) \ .withColumn("end", F.from_unixtime("end"))
最终结果
执行后即可得到目标结构的DataFrame,示例输出:
| date | id | spec_color | spec_power | spec_speed | spec_length | begin | end |
|---|---|---|---|---|---|---|---|
| 2023-08-29 | 1 | red | 155 | 198 | 4698 | 2023-08-29 | 2023-08-30 |
| 2023-08-29 | 2 | blue | 199 | 220 | 4540 | 2023-08-29 | 2023-08-30 |
内容的提问来源于stack exchange,提问作者DDR
相关产品推荐
相关产品推荐

