Spark注册临时自定义函数触发ParseException语法错误求助
问题解决:Spark SQL临时函数注册语法错误
错误原因
你的CREATE TEMPORARY FUNCTION语句不符合Spark SQL的语法规范,核心问题有两点:
- 误用了
RETURN关键字声明函数体 - 缺少
LANGUAGE SQL指定函数实现的语言类型
Spark SQL解析器无法识别这种写法,因此抛出ParseException。
修正方案
方案1:修正SQL语句语法
调整registerSQL方法中的CREATE语句,使用Spark SQL支持的标准格式:
override def registerSQL(): Unit = spark.sql( """ |CREATE OR REPLACE TEMPORARY FUNCTION lowerCaseColumn(inputValue STRING) |RETURNS STRING |LANGUAGE SQL |AS 'lower(inputValue)' |""".stripMargin)
方案2:使用Spark API直接注册UDF(更简洁)
跳过SQL语句注册,改用Spark原生的UDF注册API,彻底避免语法问题:
override def registerSQL(): Unit = { // 自定义实现转小写逻辑 spark.udf.register("lowerCaseColumn", (inputValue: String) => inputValue.toLowerCase()) // 或直接复用Spark内置的lower函数 // spark.udf.register("lowerCaseColumn", org.apache.spark.sql.functions.lower _) }
测试代码验证
修正后,测试方法中的函数调用即可正常执行。你可以添加断言来验证结果正确性:
test("should test if UDF got executed or not") { val sparkSessionBuilder = new SparkSessionBuilder import sparkSessionBuilder.spark.implicits._ val df = List("sue", "fan").toDF("input") df.createOrReplaceTempView("lake") new LowerCaseColumn(sparkSessionBuilder).registerSQL() val result = sparkSessionBuilder.spark.sql("SELECT lowerCaseColumn(input) as result FROM lake") // 添加断言验证结果 val expected = List("sue", "fan").map(_.toLowerCase).toDF("result") assert(result.collect() sameElements expected.collect()) }
内容的提问来源于stack exchange,提问作者Chandra Prakash
相关产品推荐
相关产品推荐

