You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

求含Filter条件的Databricks SQL等效PySpark代码

问题:将带FILTER条件的Databricks SQL转换为PySpark代码

我有一段包含filter条件的Databricks SQL代码,想转成PySpark代码但没思路。谷歌搜索只找到PySpark的filter基础示例,找不到适配我场景的案例。

我的Databricks SQL代码:

select
account_doc_num,

ROUND(coalesce(Sum(tbl.Amnt) FILTER(WHERE tbl.AmntType='type1' AND tbl.line_type='D'),0),2) AS Sls_Inv,
ROUND(coalesce(Sum(tbl.Amnt) FILTER(WHERE tbl.AmntType='type1' AND tbl.line_type='GVAT'),0),2) AS GVAT,

ROUND(coalesce(Sum(tbl.Amnt) FILTER(WHERE tbl.AmntType='type2' AND tbl.line_type='vat'),0),2) AS VAT,
ROUND(coalesce(Sum(tbl.Amnt) FILTER(WHERE tbl.AmntType='type2' AND tbl.line_type='K'),0),2) AS Pur_Inv,

ROUND((Sls_Inv + GVAT + VAT + Pur_Inv),2) AS Diff
    
from
DeltaTable as tbl
group by account_doc_num

我已经写了部分PySpark代码,但不知道怎么给sum函数加对应的filter条件,需要Databricks中带WHERE条件的FILTER对应的等效PySpark代码:

df_group=(
    df_vat.groupby("account_doc_num")
    .agg(
       sum("Amnt").alias("Sls_Inv"),
       sum("Amnt").alias("GVAT"),
       sum("Amnt").alias("VAT"),
       sum("Amnt").alias("Pur_Inv")
    )
)

解决方案

在PySpark中,SQL里的SUM(col) FILTER(WHERE 条件)可以用sum(when(条件, col))实现,结合coalesce处理空值、round保留小数,最后计算Diff列。完整代码如下:

from pyspark.sql import functions as F

df_group = (
    df_vat.groupBy("account_doc_num")
    .agg(
        # 对应Sls_Inv:AmntType='type1'且line_type='D'的Amnt求和
        F.round(F.coalesce(F.sum(F.when((F.col("AmntType") == "type1") & (F.col("line_type") == "D"), F.col("Amnt"))), 0), 2).alias("Sls_Inv"),
        # 对应GVAT:AmntType='type1'且line_type='GVAT'的Amnt求和
        F.round(F.coalesce(F.sum(F.when((F.col("AmntType") == "type1") & (F.col("line_type") == "GVAT"), F.col("Amnt"))), 0), 2).alias("GVAT"),
        # 对应VAT:AmntType='type2'且line_type='vat'的Amnt求和
        F.round(F.coalesce(F.sum(F.when((F.col("AmntType") == "type2") & (F.col("line_type") == "vat"), F.col("Amnt"))), 0), 2).alias("VAT"),
        # 对应Pur_Inv:AmntType='type2'且line_type='K'的Amnt求和
        F.round(F.coalesce(F.sum(F.when((F.col("AmntType") == "type2") & (F.col("line_type") == "K"), F.col("Amnt"))), 0), 2).alias("Pur_Inv")
    )
    # 计算Diff列
    .withColumn("Diff", F.round(F.col("Sls_Inv") + F.col("GVAT") + F.col("VAT") + F.col("Pur_Inv"), 2))
)

代码说明:

  • F.when(条件, 列):条件满足时返回对应列的值,否则返回null,sum时只会累加符合条件的行
  • F.coalesce(..., 0):若求和结果为null(无符合条件的行),替换为0
  • F.round(..., 2):将结果保留2位小数
  • 用withColumn计算Diff,直接引用聚合后的列即可

内容的提问来源于stack exchange,提问作者umsateesh

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.16 09:04:56