PySpark中按Type匹配Struct列键并转换为结构化数据集求助
PySpark结构化转换:从Struct类型列提取指定前缀键值对
问题背景
处理供应商提供的数据集,其中Data列为Struct类型,包含多组键值对。需要根据Type列的值,筛选Data中对应前缀的键值对,将键按下划线拆分提取班次号,转换为包含ID、Type、Shift、Role的结构化表格。
输入示例
| ID | Type | 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'} |
期望输出
| 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 |
初步思路
- 用
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
相关产品推荐
相关产品推荐

