为何SQL中array_contains支持列参数而Dataset API不行?如何解决?
问题解答:Spark中array_contains在SQL vs Dataset API的差异及解决方案
为什么Dataset API的array_contains不支持列参数?
这本质是两个API对array_contains的设计逻辑差异:
- Spark SQL里的
array_contains是作为SQL表达式解析的,它会把传入的两个参数都当作列引用(或表达式)处理,天然支持列作为第二个参数。 - 而Dataset API中的
org.apache.spark.sql.functions.array_contains方法,从你贴的报错信息就能看出来,它的第二个参数被设计成仅接受字面量(Literal),而非Column对象。当你传入$"cd"时,它会尝试把Column类型转成Literal,直接触发了类型不匹配的异常。
不用SQL表达式的解决方案
这里给你两个纯Dataset API风格的实现方式:
方式1:使用高阶函数exists(Spark 2.3+支持)
exists可以遍历数组中的每个元素,判断是否存在满足条件的元素,完美替代array_contains的列参数场景:
import org.apache.spark.sql.functions.exists val q = codes.where(exists($"codes", (elem: Column) => elem === $"cd")) q.show()
执行后会得到和SQL一致的结果:
+---------+---+ | codes| cd| +---------+---+ |[1, 2, 3]| 2| | [1]| 1| +---------+---+
方式2:使用expr函数封装SQL表达式
如果你习惯写SQL式逻辑,但不想用完整的SQL WHERE语句,可以用expr把array_contains表达式转换成Column对象:
import org.apache.spark.sql.functions.expr val q = codes.where(expr("array_contains(codes, cd)")) q.show()
这个方法本质是让Spark解析SQL表达式,但调用方式是Dataset API的风格,也符合你的需求。
额外选项:自定义UDF(不推荐)
虽然不建议(内置函数能享受Spark的优化,UDF会失去这些优势),但你也可以自定义UDF实现相同逻辑:
val arrayContainsCol = udf((arr: Seq[Int], elem: Int) => arr.contains(elem)) val q = codes.where(arrayContainsCol($"codes", $"cd")) q.show()
内容的提问来源于stack exchange,提问作者Jacek Laskowski
相关产品推荐
相关产品推荐

