Spark 3.3.2 UDF调用出现Log4j类型转换异常的解决咨询
Spark UDF日志类冲突问题解决方案
问题背景
我有一个基于IntelliJ的Scala项目(使用Spark 3.3.2),通过assembly构建胖jar包,情况如下:
- 本地命令行可正常构建并运行该胖jar;
- jar包中的代码需作为数据库UDF被调用;
- jar包可在数据库服务器上正常运行。
但调用UDF时返回如下错误:
java.lang.ClassCastException: org.apache.logging.log4j.simple.SimpleLogger incompatible with org.apache.logging.log4j.core.Logger at org.apache.spark.internal.Logging.initializeLogging(Logging.scala:130) at org.apache.spark.internal.Logging.initializeLogIfNecessary(Logging.scala:115) at org.apache.spark.internal.Logging.initializeLogIfNecessary$(Logging.scala:109) at org.apache.spark.SparkContext.initializeLogIfNecessary(SparkContext.scala:84) at org.apache.spark.internal.Logging.initializeLogIfNecessary(Logging.scala:106) at org.apache.spark.internal.Logging.initializeLogIfNecessary$(Logging.scala:105) at org.apache.spark.SparkContext.initializeLogIfNecessary(SparkContext.scala:84) at org.apache.spark.internal.Logging.log(Logging.scala:53) at org.apache.spark.internal.Logging.log$(Logging.scala:51) at org.apache.spark.SparkContext.log(SparkContext.scala:84) at org.apache.spark.internal.Logging.logInfo(Logging.scala:61) at org.apache.spark.internal.Logging.logInfo$(Logging.scala:60) at org.apache.spark.SparkContext.logInfo(SparkContext.scala:84) at org.apache.spark.SparkContext.<init>(SparkContext.scala:195) at org.apache.spark.SparkContext$.getOrCreate(SparkContext.scala:2714) at org.apache.spark.sql.SparkSession$Builder.$anonfun$getOrCreate$2(SparkSession.scala:953) at org.apache.spark.sql.SparkSession$Builder$$Lambda$111.000000007C7BCFC0.apply(Unknown Source) at scala.Option.getOrElse(Option.scala:201) at org.apache.spark.sql.SparkSession$Builder.getOrCreate(SparkSession.scala:947) at ParquetReader$.main(ParquetReader.scala:21) at ParquetReader.main(ParquetReader.scala)
经排查,问题源于作为UDF调用时,spark.internal.logging.initializeLogging初始化了与本地运行时不同的日志类(org.apache.logging.log4j.simple.SimpleLogger),现针对两个问题给出解决方案:
问题1:直接控制Spark使用的日志类
可以控制,核心是解决类加载冲突并强制指定日志实现:
1. 调整assembly打包配置
在sbt的assembly插件配置中,排除冗余的日志实现,确保只保留Log4j Core相关依赖。修改build.sbt:
assemblyMergeStrategy in assembly := { case PathList("META-INF", _*) => MergeStrategy.discard case "log4j.properties" | "log4j2.properties" => MergeStrategy.first case _ => MergeStrategy.first } libraryDependencies ++= Seq( // 匹配Spark 3.3.2对应的Log4j版本(Spark 3.3.x默认用2.17.1) "org.apache.logging.log4j" % "log4j-core" % "2.17.1" % Provided, "org.apache.logging.log4j" % "log4j-api" % "2.17.1" % Provided, "org.apache.logging.log4j" % "log4j-slf4j-impl" % "2.17.1" % Provided )
- 若数据库环境已预装Log4j Core,用
Provided避免重复打包;若环境没有,去掉Provided将依赖打包进胖jar。
2. 强制指定日志上下文工厂
在UDF代码的最开头(SparkSession初始化之前)添加系统属性设置,确保Spark加载指定的日志类:
// 强制使用Log4j Core的上下文工厂 System.setProperty("log4j2.loggerContextFactory", "org.apache.logging.log4j.core.impl.Log4jContextFactory")
问题2:通过内置配置文件完全禁用日志
可以实现,需覆盖Hadoop默认的日志配置:
1. 创建自定义日志配置文件
在项目src/main/resources目录下新建log4j2.properties,内容如下:
status = off name = DisableAllLogging # 根日志级别设为OFF,禁用所有日志输出 rootLogger.level = OFF rootLogger.appenderRefs = nullAppender rootLogger.appenderRef.nullAppender.ref = NullAppender # 定义空输出的Appender appender.nullAppender.type = NullAppender appender.nullAppender.name = NullAppender
2. 调整打包策略确保配置优先加载
在sbt的assemblyMergeStrategy中添加规则,保证自定义配置文件被优先保留:
case "log4j2.properties" => MergeStrategy.first
3. 强制加载自定义配置
在UDF代码开头添加系统属性,跳过Hadoop默认配置:
// 强制加载jar内的自定义Log4j2配置 System.setProperty("log4j.configurationFile", "log4j2.properties") // 覆盖Hadoop的根日志配置 System.setProperty("hadoop.root.logger", "OFF,nullAppender")
内容的提问来源于stack exchange,提问作者Liam385
相关产品推荐
相关产品推荐

