分步解析PySpark代码及两个join连接语句的实现作用
PySpark代码执行逻辑与Join语句作用说明
分步执行逻辑
- 构建员工表DataFrame
先定义包含4条员工记录的列表emp,对应字段依次为员工ID、姓名、上级员工ID、入职年份、所属部门ID、性别、薪资;接着指定字段名列表empColumns作为Schema,调用spark.createDataFrame()生成员工表empDF;后续调用printSchema()打印表结构、show()输出全量员工数据,truncate=False参数保证长内容不会被截断。
- 构建员工表DataFrame
- 构建部门表DataFrame
定义包含4条部门记录的列表dept,对应字段为部门名称、部门ID;指定SchemadeptColumns后调用spark.createDataFrame()生成部门表deptDF,同样打印结构和输出全量部门数据。
- 构建部门表DataFrame
- 执行关联查询并输出结果
先后执行两次关联操作,均以员工表.emp_dept_id = 部门表.dept_id为关联条件,分别使用外连接、右连接两种关联类型,输出关联后的全量结果。
- 执行关联查询并输出结果
两段Join语句的设计作用
语句1:outer全外连接
代码:empDF.join(deptDF,empDF.emp_dept_id == deptDF.dept_id,"outer")
作用:同时保留左右两张表的所有行,关联条件匹配成功的行会拼接两表字段,匹配失败的行对应另一侧表的字段填充null。
适用场景:需要全量校验两边数据的匹配情况时使用,比如本案例中可以同时查到「有对应部门的员工」、「无归属员工的空部门」、「所属部门不存在的异常员工」三类数据,不会遗漏任意一侧的原始记录。
语句2:right右连接
代码:empDF.join(deptDF,empDF.emp_dept_id == deptDF.dept_id,"right")
作用:完整保留右表(本案例中是部门表deptDF)的所有行,左表(员工表)仅保留关联条件匹配成功的行,匹配失败的行左表字段填充null。
适用场景:需要以右表为基准校验匹配情况时使用,比如本案例中可以快速筛选出没有任何员工归属的空部门(本案例中的Sales、IT两个部门),相比外连接更聚焦右表的完整性校验,不会出现左表存在但右表不存在的异常数据。
内容的提问来源于stack exchange,提问作者Sanjay KS
相关产品推荐
相关产品推荐

