如何将Python字典列表作为新列添加到PySpark DataFrame中?
问题描述
我有一个Python字典列表,希望将其作为新列添加到PySpark DataFrame中,列表中的每个元素对应新列的一个单元格。尝试过转Pandas DataFrame添加列再转回PySpark,但数据损坏;也试过用lit函数和StructType,都没成功,求解决方法。
示例数据
Python字典列表:
my_list = [ {"example1": {"subkey_example1": {"ex1": "2", "ex2": "4"}, "example2": {"subkey_example2": {"ex3": "1", "ex4": "3"}}}, {"example3": {"subkey_example3": {"ex5": "4", "ex6": "2"}, "example4": {"subkey_example4": {"ex7": "1", "ex8": "3"}}} ]
原PySpark DataFrame:
+-------+-------+ |column1|column2| +-------+-------+ | data1 | data1 | | data2 | data2 | +-------+-------+
期望结果:
+-------+-------+--------------------------------+ |column1|column2| column3 | +-------+-------+--------------------------------+ | data1 | data1 | {"example1": | | | | {"subkey_example1": | | | | {"ex1": "2", "ex2": "4"}, | | | | "example2": | | | | {"subkey_example2": | | | | {"ex3": "1", "ex4": "3"}}} | +-------+-------+--------------------------------+ | data2 | data2 | {"example3": | | | | {"subkey_example3": | | | | {"ex5": "4", "ex6": "2"}, | | | | "example4": | | | | {"subkey_example4": | | | | {"ex7": "1", "ex8": "3"}}} | +-------+-------+--------------------------------+
解决方法
方法一:通过索引列合并DataFrame
PySpark DataFrame本身没有行索引,可通过给原DataFrame和字典列表分别添加索引,再合并的方式实现需求:
- 给原DataFrame添加自增索引列
- 将字典列表转为带索引的PySpark DataFrame
- 按索引列合并两个DataFrame,最后删除索引列
代码示例(存储为字符串格式)
from pyspark.sql import SparkSession from pyspark.sql.functions import monotonically_increasing_id from pyspark.sql.types import StringType, StructType, StructField # 初始化SparkSession spark = SparkSession.builder.appName("AddDictColumn").getOrCreate() # 构造原DataFrame original_data = [("data1", "data1"), ("data2", "data2")] df = spark.createDataFrame(original_data, ["column1", "column2"]) # 给原DF添加索引 df_with_index = df.withColumn("index", monotonically_increasing_id()) # 将字典列表转为带索引的DF dict_data = [(i, str(my_list[i])) for i in range(len(my_list))] dict_schema = StructType([ StructField("index", StringType(), True), StructField("column3", StringType(), True) ]) dict_df = spark.createDataFrame(dict_data, dict_schema).withColumn("index", monotonically_increasing_id()) # 合并并删除索引列 result_df = df_with_index.join(dict_df, on="index", how="inner").drop("index") result_df.show(truncate=False)
代码示例(解析为结构化类型)
如果需要保留字典的嵌套结构而非字符串,可结合from_json解析JSON字符串:
from pyspark.sql.functions import from_json from pyspark.sql.types import StructType, StructField, StringType # 定义嵌套结构的Schema sub_schema_ex1 = StructType([ StructField("ex1", StringType()), StructField("ex2", StringType()) ]) sub_schema_ex2 = StructType([ StructField("ex3", StringType()), StructField("ex4", StringType()) ]) sub_schema_ex3 = StructType([ StructField("ex5", StringType()), StructField("ex6", StringType()) ]) sub_schema_ex4 = StructType([ StructField("ex7", StringType()), StructField("ex8", StringType()) ]) top_schema = StructType([ StructField("example1", StructType([StructField("subkey_example1", sub_schema_ex1)])), StructField("example2", StructType([StructField("subkey_example2", sub_schema_ex2)])), StructField("example3", StructType([StructField("subkey_example3", sub_schema_ex3)])), StructField("example4", StructType([StructField("subkey_example4", sub_schema_ex4)])) ]) # 构造带JSON字符串的字典DF dict_data = [(i, spark.sparkContext.parallelize([my_list[i]]).toDF().toJSON().first()) for i in range(len(my_list))] dict_df = spark.createDataFrame(dict_data, ["index", "column3_json"]).withColumn("index", monotonically_increasing_id()) # 合并并解析JSON为结构化列 result_df = df_with_index.join(dict_df, on="index", how="inner") \ .withColumn("column3", from_json("column3_json", top_schema)) \ .drop("index", "column3_json") result_df.show(truncate=False)
方法二:UDF结合行号索引
如果原DataFrame行数和字典列表长度严格一致,可通过行号索引从列表中提取对应字典:
from pyspark.sql.functions import udf, row_number from pyspark.sql.types import StringType from pyspark.sql.window import Window # 构造原DataFrame original_data = [("data1", "data1"), ("data2", "data2")] df = spark.createDataFrame(original_data, ["column1", "column2"]) # 生成连续行号索引 window = Window.orderBy("column1") df_with_row_num = df.withColumn("row_num", row_number().over(window)) # 定义UDF根据行号取字典 @udf(StringType()) def get_dict(row_num): return str(my_list[row_num - 1]) # row_number从1开始,列表索引从0开始 # 添加新列并删除行号 result_df = df_with_row_num.withColumn("column3", get_dict("row_num")).drop("row_num") result_df.show(truncate=False)
内容的提问来源于stack exchange,提问作者duck17
相关产品推荐
相关产品推荐

