Spark SQL调用from_avro函数报错,求注册与正确调用方法
解决Spark SQL中
from_avro函数未定义的问题 1. 确保引入Spark Avro依赖
from_avro属于Spark的Avro扩展库,默认不包含在核心Spark包中,必须先添加对应版本的依赖:
Spark Submit命令行:
spark-submit \ --packages org.apache.spark:spark-avro_2.12:3.5.0 \ your-application.jar注意替换
2.12为你的Scala版本,3.5.0为匹配的Spark版本。SBT项目:
在build.sbt中添加:libraryDependencies += "org.apache.spark" %% "spark-avro" % "3.5.0" % ProvidedMaven项目:
在pom.xml中添加:<dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-avro_2.12</artifactId> <version>3.5.0</version> <scope>provided</scope> </dependency>
2. 注册from_avro函数到Spark SQL
添加依赖后,需显式将函数注册到Spark SQL的函数注册表,才能在SQL语句中直接调用:
方法一:注册临时函数
在代码中添加以下行,注册后即可在SQL中使用:
import org.apache.spark.sql.avro.functions.from_avro ss.udf.register("from_avro", from_avro _)
这里的ss是你的SparkSession实例。
方法二:通过DataFrame API结合临时视图
如果不想注册函数,也可以先通过DataFrame API处理,再转为SQL可用的视图:
import org.apache.spark.sql.avro.functions.from_avro val processedDf = testDF.select(from_avro($"avro_payload", $"avro_schema_string").alias("avro_data")) processedDf.createOrReplaceTempView("processedTestDF") // 后续可在SQL中查询处理后的数据 val rsDf = ss.sql("SELECT avro_data FROM processedTestDF")
3. 验证调用
完成上述步骤后,原SQL语句即可正常执行:
val rsDf = ss.sql("SELECT from_avro(avro_payload, avro_schema_string) FROM testDF")
内容的提问来源于stack exchange,提问作者Philip K. Adetiloye
相关产品推荐
相关产品推荐

