PySpark对比数据集识别未变更记录:anti join空值判断问题
回答
关于null值相等的join配置
Spark 2.3.0及以上版本提供了参数spark.sql.join.nullSafeEquality,设置为true时,所有join的等值匹配逻辑会将两侧同列的null值判定为相等,完全符合你要的效果。
可以在提交作业时配置,也可以在代码里动态设置:
spark.conf.set("spark.sql.join.nullSafeEquality", "true")
配置生效后,你最开始写的left_anti join代码不需要任何修改,就能直接得到正确结果,不会把双null的记录误判为变更。
注意:这个配置仅修改join场景下的等值比较逻辑,不会影响where、select等其他位置的三值逻辑判断。如果你的作业里存在其他join逻辑依赖null不相等的规则,不建议开全局配置,避免逻辑错误。
无需修改全局配置的更优实现
如果不想改全局配置,直接用Spark内置的*空值安全等号运算符<=>*写join条件就行,比你现在用的exceptAll+回连的方案更简洁,性能也更好,只需要一次join shuffle,不需要额外回表关联技术元数据列。
代码实现:
# 主键+业务列都用<=>做关联,自动兼容null值相等判断 null_safe_join_cond = (new["text"] <=> old["text"]) & (new["first_character"] <=> old["first_character"]) changed_records = new.join(old, null_safe_join_cond, "left_anti") changed_records.show()
运行后输出正好是你预期的结果,仅返回真正发生变更的Bar记录:
+----+---------------+-------+ |text|first_character|version| +----+---------------+-------+ | Bar| B| 2| +----+---------------+-------+
方案对比
- 你当前的
exceptAll方案:需要先对主键+业务列做一次全量shuffle比对差异,再二次join关联回version这类技术列,多了一次shuffle开销,业务列多的时候要重复写两遍列名,维护成本高。 - 空值安全等号方案:仅需一次left_anti join就完成变更识别,技术元数据列会自动保留,不需要额外select和join,非常适合封装成通用的历史化处理函数——只要传入主键列列表、业务比对列列表,就能自动拼接join条件,自动兼容null值判断。
- 全局配置方案:适合整个ETL作业所有join都需要null判等的场景,代码量最少,但要提前确认作业内其他join逻辑不会受该规则影响。
内容的提问来源于stack exchange,提问作者Alfred G
相关产品推荐
相关产品推荐

