将可用SQL关联查询转换为PySpark DataFrame语法的技术咨询
把验证通过的SQL查询转成PySpark DataFrame语法
嘿,我来帮你把这段已经验证可用的SQL转换成PySpark的DataFrame操作写法,咱们先拆解原SQL的逻辑,再对应实现:
原SQL(查找存在匹配record_id的“无效”记录)
SELECT DISTINCT(r.id) FROM core_record AS r LEFT OUTER JOIN core_recordvalidation rv ON r.id = rv.record_id WHERE r.job_id = 41 AND rv.record_id is not null;
对应的PySpark代码实现
假设你已经加载好了两个DataFrame:core_record_df(对应SQL中的core_record表)和core_recordvalidation_df(对应core_recordvalidation表),代码如下:
# 先给DataFrame起别名,和SQL里的AS r/AS rv对应,方便后续引用列 r_df = core_record_df.alias("r") rv_df = core_recordvalidation_df.alias("rv") # 执行左外连接,关联条件为r.id = rv.record_id joined_df = r_df.join( rv_df, on=r_df["id"] == rv_df["record_id"], how="left_outer" ) # 过滤符合条件的数据,去重后选取r.id列 result_df = joined_df \ .filter((r_df["job_id"] == 41) & (rv_df["record_id"].isNotNull())) \ .select(r_df["id"]) \ .distinct() # 查看结果(可选) result_df.show()
关键细节说明
- 用
alias()给DataFrame设置别名,和SQL的表别名逻辑完全一致,避免列名冲突 join()方法的how="left_outer"严格对应SQL的LEFT OUTER JOIN- 过滤条件可以用
filter()或者where(),PySpark里这两个方法功能相同,多条件用&连接时记得加括号,避免优先级错误 isNotNull()是PySpark中对应SQLis not null的内置方法distinct()实现SQL中SELECT DISTINCT的去重效果
如果后续需要选取更多列,直接在select()里添加对应的列即可,比如.select(r_df["id"], r_df["job_id"], rv_df["status"])。
内容的提问来源于stack exchange,提问作者ghukill
相关产品推荐
相关产品推荐

