如何实现含列表类型键列的两个DataFrame关联?
实现思路与代码示例
要解决这个关联问题,核心思路是先把第二个DataFrame中包含tmpObj的列表列“扁平化”——也就是把列表里的每个tmpObj拆成单独的行,这样就能用常规的DataFrame关联操作来匹配id了。下面分Spark(PySpark)和Pandas两种常用场景给你具体实现方案:
一、PySpark 实现方案
步骤1:准备示例数据
先构造符合你描述的两个DataFrame,其中第二个DF的列表列是包含struct类型元素的数组:
from pyspark.sql import SparkSession from pyspark.sql.functions import explode, col from pyspark.sql.types import StructType, StructField, StringType, IntegerType spark = SparkSession.builder.appName("ListJoinDemo").getOrCreate() # 第一个DF:id为String类型,附带其他字段 df1 = spark.createDataFrame([ ("1", "Alice"), ("2", "Bob"), ("3", "Charlie") ], ["id", "name"]) # 定义tmpObj的结构 tmp_obj_schema = StructType([ StructField("id", StringType(), True), StructField("value", IntegerType(), True) ]) # 第二个DF:包含列表类型的obj_list列,元素是tmpObj df2 = spark.createDataFrame([ ("X", [{"id": "1", "value": 10}, {"id": "2", "value": 20}]), ("Y", [{"id": "2", "value": 30}, {"id": "3", "value": 40}]) ], ["category", "obj_list"])
步骤2:扁平化列表列
用explode函数把列表中的每个tmpObj拆成单独的行(如果要保留空列表的行,可以用explode_outer):
# 展开列表,生成包含单个tmpObj的列 df2_exploded = df2.withColumn("tmp_obj", explode(col("obj_list")))
步骤3:提取tmpObj中的id和value
把struct类型的tmp_obj字段拆成单独的obj_id和value列,方便后续关联:
df2_flattened = df2_exploded.select( col("category"), col("tmp_obj.id").alias("obj_id"), col("tmp_obj.value").alias("value") )
步骤4:关联两个DataFrame
用常规的join操作,以df1.id和df2_flattened.obj_id为关联键,可根据需求选择关联类型(inner/left/right/full):
# 这里用inner join,你可以根据业务调整为left join等 joined_df = df1.join(df2_flattened, df1.id == df2_flattened.obj_id, "inner") # 查看结果 joined_df.show()
执行结果会把匹配的记录关联起来,比如id=2的Bob会关联到两条来自df2的记录。
二、Pandas 实现方案
如果是用Pandas处理,思路完全一致,只是API略有不同:
步骤1:准备示例数据
import pandas as pd df1 = pd.DataFrame({ "id": ["1", "2", "3"], "name": ["Alice", "Bob", "Charlie"] }) df2 = pd.DataFrame({ "category": ["X", "Y"], "obj_list": [ [{"id": "1", "value": 10}, {"id": "2", "value": 20}], [{"id": "2", "value": 30}, {"id": "3", "value": 40}] ] })
步骤2:扁平化列表列
用Pandas的explode方法展开列表:
df2_exploded = df2.explode("obj_list", ignore_index=True)
步骤3:提取tmpObj中的id和value
通过apply或pd.json_normalize提取字段:
# 方式1:用apply提取 df2_flattened = df2_exploded.assign( obj_id=df2_exploded["obj_list"].apply(lambda x: x["id"]), value=df2_exploded["obj_list"].apply(lambda x: x["value"]) ).drop("obj_list", axis=1) # 方式2:用json_normalize更高效(适合大数据集) # from pandas import json_normalize # df2_flattened = pd.concat([df2_exploded["category"], json_normalize(df2_exploded["obj_list"])], axis=1).rename(columns={"id": "obj_id"})
步骤4:关联两个DataFrame
用pd.merge完成关联:
joined_df = pd.merge(df1, df2_flattened, left_on="id", right_on="obj_id", how="inner") print(joined_df)
关键注意点
- 如果列表中存在空值或空列表,
explode(Spark/Pandas)会默认丢弃对应的行;若需保留,Spark用explode_outer,Pandas则会生成NaN值,可后续处理。 - 关联类型根据业务需求选择:比如想保留df1的所有记录,用
left join;想保留双方所有记录,用full join。
内容的提问来源于stack exchange,提问作者Ben Jen
相关产品推荐
相关产品推荐

