在pyspark中加载sql之后,会经常遇到各DataFrame之间的join操作,以下给出集中join的调用方式。

sql加载

该部分主要是解决初学朋友经常私信的问题,如果已经知道,可以直接跳过该部分

from pyspark.sql import SparkSession
from pyspark import SparkConf

# 加载本地数据
sparkConf = SparkConf().set("spark.driver.host", "localhost")
spark = SparkSession.builder.config(conf=sparkConf).enableHiveSupport().getOrCreate()
df = spark.read.cvs("file_path", spe="\t", header=True)


# 加载sql数据
spark = SparkSession.builder.enableHiveSupport().getOrCreate()
df = 
spark.sql(your_sql)

join 操作

整个思路,定义join的维度,依据inner、leftanti、left来确定join的方式

# 定义join维度
con = [df_1.col_1 == df_2.col_2]
 # 取两个DataFrame的inner join数据
df_1 = df_1.join(df_2, on=con, how='inner') 
 # 取两个DataFrame的left join数据
df_1 = df_1.join(df_2, on=con, how='left') 
 # 取两个DataFrame的“差集”join数据
df_1 = df_1.join(df_2, on=con, how='leftanti') 
Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐