PySpark使用when/otherwise处理空值时的语法错误排查
解决Spark数组空值拼接的语法错误及实现方案
错误原因
你的SQL表达式写法违反了Spark SQL语法规范:
x.altitude.isNotNull()是PySpark DataFrame API的方法调用,Spark SQL中判断字段非空必须用SQL原生语法x.altitude IS NOT NULL- 原写法中
when(...).otherwise(...)的链式调用不符合Spark SQL表达式的语法要求,导致解析器报错
正确实现方式
方案1:修正Spark SQL表达式
使用Spark SQL内置的IF函数或CASE WHEN语句处理空值,语法更直观:
from pyspark.sql import functions as F def stringify_litetrack_points(points): # 用IF函数替换空值为'*' return F.expr("array_join(transform(points, x -> IF(x.altitude IS NOT NULL, round(x.altitude), '*')), ':')")
或者使用CASE WHEN(逻辑更清晰,适合复杂条件):
def stringify_litetrack_points(points): return F.expr("array_join(transform(points, x -> CASE WHEN x.altitude IS NOT NULL THEN round(x.altitude) ELSE '*' END), ':')")
方案2:纯PySpark DataFrame API实现(推荐)
完全避开SQL表达式的语法陷阱,直接用DataFrame API操作,可读性和可维护性更高:
from pyspark.sql import functions as F def stringify_litetrack_points(points_col): # 遍历数组每个元素,将空altitude替换为'*' transformed = F.transform( points_col, lambda p: F.when(p["altitude"].isNotNull(), F.round(p["altitude"])) .otherwise(F.lit('*')) ) # 拼接数组为字符串 return F.array_join(transformed, ':')
调用方式保持不变:
df.withColumn("altitude_list", stringify_litetrack_points("points")).drop("points").show()
两种方案都能保留原数组的元素数量,空altitude会被替换为*后参与拼接。
内容的提问来源于stack exchange,提问作者4mla1fn
相关产品推荐
相关产品推荐

