Spark2中无法设置StreamingLogLevel及继承Logging对象求助
解决Spark 2.x中无法继承Logging及设置Streaming日志级别的问题
嘿,我完全懂你遇到的这两个头疼问题!Spark 2.x里确实把Logging trait设为私有了,而且Streaming相关的日志设置也和早期版本有点不一样,给你两个问题的具体解决办法:
一、替代私有Logging trait的方案
Spark 2.x中org.apache.spark.Logging是private[spark]修饰的,外部代码没法直接继承。不过别慌,Spark底层用的是SLF4J日志框架,咱们直接用它来实现日志功能就行,这也是官方推荐的方式:
import org.slf4j.LoggerFactory // 比如你的流处理类 class MyStreamingJob { // 获取当前类的Logger实例 private val logger = LoggerFactory.getLogger(getClass) def startJob(): Unit = { // 正常使用日志方法 logger.info("开始启动流处理任务啦") logger.debug("调试信息:当前处理的批次ID是xxx") logger.warn("注意:当前数据量超出阈值") } }
这样写完全能达到和继承Logging一样的日志效果,而且兼容性更好,后续升级Spark版本也不用改这块代码。
二、设置Spark Streaming的日志级别
你没法直接用StreamingLogLevel的话,可以通过以下几种方式调整Streaming相关的日志输出:
1. 代码中动态设置(适合临时调整)
通过SLF4J的API直接修改指定包的日志级别,比如把Streaming的日志设为WARN级别,减少冗余输出:
import org.slf4j.LoggerFactory import ch.qos.logback.classic.Level // 注意:这里需要确保你的日志框架是logback(Spark默认用的就是它) val streamingLogger = LoggerFactory.getLogger("org.apache.spark.streaming") .asInstanceOf[ch.qos.logback.classic.Logger] streamingLogger.setLevel(Level.WARN)
2. 通过SparkConf配置参数
在创建SparkConf的时候直接指定日志相关参数:
import org.apache.spark.SparkConf import org.apache.spark.streaming.StreamingContext import org.apache.spark.streaming.Seconds val conf = new SparkConf() .setAppName("MyStreamingApp") .setMaster("local[*]") // 设置Streaming模块的日志级别 .set("org.apache.spark.streaming.level", "WARN") val ssc = new StreamingContext(conf, Seconds(5))
3. 用配置文件永久设置(推荐生产环境)
在项目的src/main/resources目录下添加logback.xml(Spark默认日志框架是logback),在里面配置Streaming的日志级别:
<configuration> <!-- 根日志级别 --> <root level="INFO"> <appender-ref ref="CONSOLE"/> </root> <!-- 单独设置Spark Streaming的日志级别 --> <logger name="org.apache.spark.streaming" level="WARN"/> <!-- 如果还想调整其他Spark模块的日志,比如Core、SQL,也可以加对应的logger --> <logger name="org.apache.spark.core" level="INFO"/> </configuration>
这样就能完美解决你遇到的两个问题啦!
内容的提问来源于stack exchange,提问作者J.F
相关产品推荐
相关产品推荐

