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

PySpark中按Type匹配Struct列键并转换为结构化数据集求助

PySpark结构化转换:从Struct类型列提取指定前缀键值对

问题背景

处理供应商提供的数据集,其中Data列为Struct类型,包含多组键值对。需要根据Type列的值,筛选Data中对应前缀的键值对,将键按下划线拆分提取班次号,转换为包含ID、Type、Shift、Role的结构化表格。

输入示例

IDTypeData
1Day{Day_1:'Monitor', Day_2:'Manage', Evening_6:'None', Evening_3:'Monitor'}
2Evening{Day_2:'N/A', Day_3:'None', Evening_1:'On Call', Evening_2:'On Call'}
3Day{Day_1:'Manage', Day_5:'On Call', Evening_2:'None'}

期望输出

IDTypeShiftRole
1Day1Monitor
1Day2Manage
2Evening1On Call
2Evening2On Call
3Day1Manage
3Day5On Call

初步思路

  • 用Type列值+_+数字构建正则表达式
  • 从Struct类型的Data列中筛选符合前缀的键值对
  • 用explode函数将筛选后的键值对展开为行
  • 拆分键名的下划线部分,提取班次号

具体实现步骤与代码

1. 初始化环境与加载数据

先创建SparkSession并加载模拟数据(实际场景替换为你的数据源):

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, explode, split, regexp_extract, struct, lit, array

# 初始化SparkSession
spark = SparkSession.builder.appName("StructTransform").getOrCreate()

# 模拟输入数据(实际读取你自己的数据集即可)
raw_data = [
    (1, "Day", {"Day_1": "Monitor", "Day_2": "Manage", "Evening_6": "None", "Evening_3": "Monitor"}),
    (2, "Evening", {"Day_2": "N/A", "Day_3": "None", "Evening_1": "On Call", "Evening_2": "On Call"}),
    (3, "Day", {"Day_1": "Manage", "Day_5": "On Call", "Evening_2": "None"})
]

df = spark.createDataFrame(raw_data, ["ID", "Type", "Data"])

2. 展开Struct列的所有键值对

Struct类型无法直接按前缀筛选,先把所有字段转为key-value结构的数组,再用explode拆分为行:

# 获取Data列的所有字段名
data_fields = df.select("Data.*").columns

# 生成键值对数组并展开
df_exploded = df.withColumn(
    "key_value",
    explode(array(*[struct(lit(field).alias("key"), col(f"Data.{field}").alias("value")) for field in data_fields]))
).select("ID", "Type", "key_value.key", "key_value.value")

3. 筛选符合Type前缀的有效键值对

用正则匹配Type_数字格式的键,同时过滤掉None/N/A这类无效角色值:

df_filtered = df_exploded.filter(
    col("key").rlike(col("Type") + "_\\d+")
).filter(~col("value").isin("None", "N/A"))

4. 提取班次号并生成最终表格

拆分键名的下划线部分,提取班次号,并重命名列得到目标结构:

df_final = df_filtered.withColumn(
    "Shift",
    split(col("key"), "_")[1].cast("int")
).select("ID", "Type", "Shift", col("value").alias("Role"))

# 查看结果
df_final.show()

运行输出

+---+-------+-----+---------+
| ID|   Type|Shift|     Role|
+---+-------+-----+---------+
|  1|    Day|    1|  Monitor|
|  1|    Day|    2|   Manage|
|  2|Evening|    1| On Call|
|  2|Evening|    2| On Call|
|  3|    Day|    1|   Manage|
|  3|    Day|    5| On Call|
+---+-------+-----+---------+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 10:55:21