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

如何在示例中连接DataFrame?需为关联字段添加_src/_dst前缀

这问题我碰到过好多次,其实核心思路就是分别关联源节点和目标节点的数据,提前给关联后的列加上指定前缀,下面分Pandas和PySpark两种常用场景给你具体实现代码:

Pandas 实现方案

先把你的示例数据构造出来(方便你直接测试):

import pandas as pd

edges = pd.DataFrame({
    'srcId': [1,1,1,2,4],
    'dstId': [3,4,2,3,3],
    'timestamp': [1345534569, 1346564657, 1345769687, 1345769687, 1345769687]
})

vertices = pd.DataFrame({
    'id': [1,2,3,4],
    'name': ['abc', 'def', 'rtf', 'wrr'],
    's_type': ['A', 'B', 'C', 'D']
})

接下来分两步关联并处理列名:

  1. 关联源节点信息,提前重命名列加上_src前缀
# 关联源节点数据,同时重命名name和s_type列
result = edges.merge(
    vertices.rename(columns={'name': 'name_src', 's_type': 's_type_src'}),
    left_on='srcId',
    right_on='id',
    how='left'
).drop('id', axis=1)  # 删掉关联后多余的id列
  1. 关联目标节点信息,同样重命名列加上_dst前缀
# 关联目标节点数据,重命名列
result = result.merge(
    vertices.rename(columns={'name': 'name_dst', 's_type': 's_type_dst'}),
    left_on='dstId',
    right_on='id',
    how='left'
).drop('id', axis=1)

最后调整列的顺序,和你想要的结构完全对齐:

result = result[['srcId', 'name_src', 's_type_src', 'dstId', 'name_dst', 's_type_dst', 'timestamp']]

运行后你就能得到完全符合要求的DataFrame了。

PySpark 实现方案

如果是用Spark处理大数据,思路类似,只是语法稍有不同:
先构造Spark DataFrame:

from pyspark.sql import SparkSession
from pyspark.sql.functions import col

spark = SparkSession.builder.appName("vertex_edge_join").getOrCreate()

edges_df = spark.createDataFrame([
    (1,3,1345534569),
    (1,4,1346564657),
    (1,2,1345769687),
    (2,3,1345769687),
    (4,3,1345769687)
], ['srcId', 'dstId', 'timestamp'])

vertices_df = spark.createDataFrame([
    (1,'abc','A'),
    (2,'def','B'),
    (3,'rtf','C'),
    (4,'wrr','D')
], ['id', 'name', 's_type'])

然后通过两次join实现,每次join前先把vertices的列重命名好:

# 第一次join:关联源节点,重命名列加_src前缀
result_df = edges_df.join(
    vertices_df.select(
        col('id').alias('srcId'),
        col('name').alias('name_src'),
        col('s_type').alias('s_type_src')
    ),
    on='srcId',
    how='left'
)

# 第二次join:关联目标节点,重命名列加_dst前缀
result_df = result_df.join(
    vertices_df.select(
        col('id').alias('dstId'),
        col('name').alias('name_dst'),
        col('s_type').alias('s_type_dst')
    ),
    on='dstId',
    how='left'
)

# 调整列顺序到目标格式
result_df = result_df.select('srcId', 'name_src', 's_type_src', 'dstId', 'name_dst', 's_type_dst', 'timestamp')

这样就能得到你想要的结构,不管是小数据用Pandas还是大数据用Spark,这个思路都适用。

内容的提问来源于stack exchange,提问作者ScalaBoy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 04:03:14