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

PySpark 3.0.3中如何用expr重构aggregate函数代码?

场景1:结构体数组字段求和

你之前的错误是用f.lit()包裹了expr(),这会把SQL表达式的结果当成字符串常量,而非执行计算。另外初始值要和数值类型匹配,改用0.0。正确写法:

df.withColumn(
    "total_amount",
    expr("aggregate(taxes, 0.0, (acc, x) -> acc + x.amount)")
)
  • 不需要把taxes改成数字列表,只要它是包含amount字段的结构体数组,SQL的aggregate就能直接处理。
  • SQL里访问结构体字段也可以用x['amount'],但字段名无特殊字符时x.amount更简洁。
场景2:过滤数组后求和

直接在SQL表达式里嵌套filter和aggregate函数,用SQL语法写过滤条件即可:

df.withColumn(
    "total_amount_2",
    expr("""
        aggregate(
            filter(lines, x -> x.id IS NOT NULL AND x.code NOT IN ('CODE1', 'CODE2')),
            0.0,
            (acc, x) -> acc + x.amount
        )
    """)
)
  • filter函数负责筛选数组中id非空且code不在指定列表的元素,再将过滤后的数组传入aggregate求和。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 01:10:33