PySpark解析表头定义列名的JSON文件:列值映射难题
问题
我正尝试使用PySpark解析表头仅定义一次列名的JSON文件,目前已成功加载列名数组与数据值数组,但卡在如何将values数组的第i个值扩展为以colnames数组第i个值为别名的新列。实际数据集包含约100列。
示例JSON数据
json_example = """{ "Header": { "Foo": 'bar', "SignalList": [ {"Name": "id", "Type": "integer"}, {"Name": "field01", "Type": "float"}, {"Name": "field02", "Type": "float"} ] }, "Payload": [ { "Data": [ [10001, 2.2, -102.4], [10002, 2.3, -102.3], [10003, 2.2, -102.4], [10004, 2.6, -102.4], [10005, 2.5, -102.5] ] }, { "Data":[ [10006, 2.4, -102.3], [10007, 2.5, -102.5], [10008, 2.5, -102.5], [10009, 2.6, -102.4] ] } ] }"""
目标格式
+-------+-------+-------+------+ | id |field01|field02| foo| +-------+-------+-------+------+ | 10001| 2.2| -102.4| bar| | 10002| 2.3| -102.3| bar| | 10003| 2.2| -102.4| bar| | 10004| 2.6| -102.4| bar| | 10005| 2.5| -102.5| bar| | 10006| 2.4| -102.3| bar| | 10007| 2.5| -102.5| bar| | 10008| 2.5| -102.5| bar| | 10009| 2.6| -102.4| bar| +-------+-------+-------+------+
当前已完成的代码
# Read JSON df = spark.read.option("multiLine","true").json(sc.parallelize([json_example])) # Copy df1 = df # Get column names df1 = df1.withColumn("colnames", F.col('Header.SignalList.Name')) # Get data values df1 = df1.withColumn("values", F.explode('Payload.Data').alias('values')) df1 = df1.withColumn("values", F.explode('values')) # Get metadata from Header df1 = df1.withColumn("Foo", F.col("Header.Foo")) df1 = df1.select(["foo","colnames", "values"]) df1.show(truncate=False)
当前输出结果
+---+----------------------+----------------------+ |foo|colnames |values | +---+----------------------+----------------------+ |bar|[id, field01, field02]|[10001.0, 2.2, -102.4]| |bar|[id, field01, field02]|[10002.0, 2.3, -102.3]| |bar|[id, field01, field02]|[10003.0, 2.2, -102.4]| |bar|[id, field01, field02]|[10004.0, 2.6, -102.4]| |bar|[id, field01, field02]|[10005.0, 2.5, -102.5]| |bar|[id, field01, field02]|[10006.0, 2.4, -102.3]| |bar|[id, field01, field02]|[10007.0, 2.5, -102.5]| |bar|[id, field01, field02]|[10008.0, 2.5, -102.5]| |bar|[id, field01, field02]|[10009.0, 2.6, -102.4]| +---+----------------------+----------------------+
解决方案
要将values数组元素映射为colnames对应列名的字段,可按以下步骤操作:
1. 提取统一列名列表
由于整个数据集的列名是全局定义的,直接取第一行的colnames数组值即可:
col_names = df1.select("colnames").first()[0]
2. 动态生成列映射规则
使用PySpark的element_at函数(索引从1开始),遍历列名列表,将values数组对应位置的元素取出并命名为目标列名:
from pyspark.sql import functions as F # 生成列表达式集合 columns = [ F.element_at(F.col("values"), i+1).alias(col_names[i]) for i in range(len(col_names)) ] # 保留原有的foo列 columns.append(F.col("foo"))
3. 重构得到目标DataFrame
通过select方法应用上述列表达式,即可得到符合要求的结果:
final_df = df1.select(*columns) final_df.show()
完整可运行代码
from pyspark.sql import functions as F # 读取JSON数据 df = spark.read.option("multiLine","true").json(sc.parallelize([json_example])) df1 = df # 提取列名数组 df1 = df1.withColumn("colnames", F.col('Header.SignalList.Name')) # 展开Payload中的数据数组 df1 = df1.withColumn("values", F.explode('Payload.Data')) df1 = df1.withColumn("values", F.explode('values')) # 提取Header中的Foo字段 df1 = df1.withColumn("foo", F.col("Header.Foo")) df1 = df1.select(["foo","colnames", "values"]) # 获取列名列表 col_names = df1.select("colnames").first()[0] # 动态生成列映射 columns = [F.element_at(F.col("values"), i+1).alias(col_names[i]) for i in range(len(col_names))] columns.append(F.col("foo")) # 构建最终DataFrame final_df = df1.select(*columns) # 可选:根据SignalList中的Type字段转换数据类型 for field in df.select("Header.SignalList").first()[0]: final_df = final_df.withColumn( field["Name"], F.col(field["Name"]).cast(field["Type"]) ) final_df.show()
关键说明
- 这种动态生成列的方式完全适配100列的大场景,无需手动编写重复代码。
- 可选的类型转换步骤能让结果数据类型与定义的
Type完全匹配,提升数据准确性。
内容的提问来源于stack exchange,提问作者user3311719
相关产品推荐
相关产品推荐

