如何将指定SQL查询转换为Databricks DataFrame的selectExpr格式构建查询
将指定SQL转换为Databricks PySpark DataFrame代码(使用selectExpr)
以下是对应原SQL查询的PySpark实现,通过DataFrame API和selectExpr完成:
步骤1:构建各子查询对应的DataFrame
处理table2的聚合逻辑
from pyspark.sql import functions as F # 生成table2的聚合结果(对应原SQL中的mm子查询) mm_df = spark.table("table2") \ .groupBy("ppcolumn1") \ .agg(F.max("ppcolumn1").alias("max_ppcolumn1"))
构建a子查询的DataFrame
# 关联table1、mm_df和table3,得到a_df a_df = spark.table("table1").alias("ss") \ .join(mm_df.alias("mm"), (F.col("ss.ppcolumn1") == F.col("mm.ppcolumn1")) & (F.col("ss.ppcolumn1") == F.col("mm.max_ppcolumn1")), "inner") \ .join(spark.table("table3").alias("Ph"), F.col("ss.columnid") == F.col("Ph.columnid"), "inner") \ .selectExpr("ss.*", "Ph.testcol AS column1")
构建b子查询的DataFrame
# 生成table3的聚合结果(对应原SQL中的mmsq子查询) mmsq_b = spark.table("table3") \ .groupBy("cscol") \ .agg(F.max("csacol1").alias("max_csacol")) # 关联得到b_df b_df = spark.table("table3").alias("ssco") \ .join(mmsq_b.alias("mmsq"), (F.col("ssco.cscol") == F.col("mmsq.cscol")) & (F.col("ssco.csacol1") == F.col("mmsq.max_csacol")), "inner") \ .selectExpr("ssco.*")
构建c子查询的DataFrame
# 生成table4的聚合结果(对应原SQL中的mmsq子查询) mmsq_c = spark.table("table4") \ .groupBy("adcol") \ .agg(F.max("colsq").alias("max_colsq")) # 关联得到c_df c_df = spark.table("table4").alias("ssad") \ .join(mmsq_c.alias("mmsq"), (F.col("ssad.adcol") == F.col("mmsq.adcol")) & (F.col("ssad.colsq") == F.col("mmsq.max_colsq")), "inner") \ .selectExpr("ssad.*")
步骤2:关联所有DataFrame并选择目标列
# 左连接a_df、b_df、c_df,最终筛选所需字段 final_df = a_df.alias("a") \ .join(b_df.alias("b"), F.col("a.colid") == F.col("b.cscol"), "left_outer") \ .join(c_df.alias("c"), F.col("a.colid") == F.col("c.adcol"), "left_outer") \ .selectExpr( "a.column1", "b.ccname AS columnname", "c.aname AS addcolumn1", "c.aaname AS addcolumn2", "c.tname AS columnc", "c.ssname AS columns", "c.ccname AS columncc" )
说明
- 所有子查询拆分实现,避免嵌套SQL的复杂度
selectExpr直接映射原SQL中的字段别名与选择逻辑- 连接条件严格匹配原SQL的ON子句规则
内容的提问来源于stack exchange,提问作者Developer Rajinikanth
相关产品推荐
相关产品推荐

