You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何实现含列表类型键列的两个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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.20 09:18:12