如何在示例中连接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'] })
接下来分两步关联并处理列名:
- 关联源节点信息,提前重命名列加上
_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列
- 关联目标节点信息,同样重命名列加上
_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
相关产品推荐
相关产品推荐

