如何基于UPDATED_FIELDS列转换Spark DataFrame结构?
Spark DataFrame 拆分数组列并映射字段值
解决方案思路
先将UPDATED_FIELDS数组列拆分为多行,每个数组元素对应单独一行;再根据拆分出的field值,匹配提取s前缀和d前缀对应列的数值。
代码实现(PySpark)
1. 构造示例DataFrame
from pyspark.sql import SparkSession from pyspark.sql.functions import explode, col, when, create_map, lit spark = SparkSession.builder.appName("field_extract").getOrCreate() # 原始数据 data = [ (1, "AAA", "BBB", 10, "AAA__", "BBB", 10, ["first"]), (2, "CCC", "DDD", 20, "CCC__", "DDD", 21, ["first", "age"]) ] raw_df = spark.createDataFrame( data, ["s.id", "s.first", "s.last", "s.age", "d.first", "d.last", "d.age", "UPDATED_FIELDS"] )
2. 拆分数组并映射字段
方式一:条件判断(适合字段较少场景)
# 重命名ID列 processed_df = raw_df.withColumnRenamed("s.id", "id") # 拆分数组,每行对应一个field exploded_df = processed_df.select( "id", explode("UPDATED_FIELDS").alias("field"), "s.first", "s.age", "d.first", "d.age" ) # 根据field匹配对应值 result_df = exploded_df.select( "id", "field", when(col("field") == "first", col("s.first")) .when(col("field") == "age", col("s.age")) .alias("s_value"), when(col("field") == "first", col("d.first")) .when(col("field") == "age", col("d.age")) .alias("d_value") ).drop("s.first", "s.age", "d.first", "d.age") # 查看结果 result_df.show()
方式二:映射表匹配(适合字段较多场景)
# 重命名ID列并拆分数组 processed_df = raw_df.withColumnRenamed("s.id", "id") exploded_df = processed_df.select("id", explode("UPDATED_FIELDS").alias("field"), "*") # 创建字段映射关系 s_field_map = create_map(lit("first"), col("s.first"), lit("age"), col("s.age")) d_field_map = create_map(lit("first"), col("d.first"), lit("age"), col("d.age")) # 提取对应值 result_df = exploded_df.select( "id", "field", s_field_map.getItem(col("field")).alias("s_value"), d_field_map.getItem(col("field")).alias("d_value") ).drop(*[col for col in raw_df.columns if col != "id"]) # 查看结果 result_df.show()
最终输出
+---+-----+-------+-------+ | id|field|s_value|d_value| +---+-----+-------+-------+ | 1|first| AAA| AAA__| | 2|first| CCC| CCC__| | 2| age| 20| 21| +---+-----+-------+-------+
内容的提问来源于stack exchange,提问作者gherkin
相关产品推荐
相关产品推荐

