PySpark合并小文件时[empty row]日志消息的屏蔽方法求助
屏蔽PySpark FileScanRDD的"partition values: [empty row]"日志
你可以通过调整FileScanRDD类的日志级别来屏蔽这类INFO级别的消息,以下是几种可行的实现方式:
方法1:在PySpark代码内动态设置日志级别
在创建SparkSession后,添加日志配置代码,直接将FileScanRDD的日志级别设为WARN或ERROR(高于INFO的级别即可):
import sys from pyspark.sql import SparkSession from py4j.java_gateway import java_import if len(sys.argv) != 3: raise Exception("bad command line arguments %s" % (sys.argv) ) input_directory = sys.argv[1] output_directory = sys.argv[2] spark = SparkSession.builder.getOrCreate() # 配置FileScanRDD日志级别 java_import(spark._jvm, "org.apache.log4j.Logger") java_import(spark._jvm, "org.apache.log4j.Level") logger = spark._jvm.Logger.getLogger("org.apache.spark.sql.execution.datasources.FileScanRDD") logger.setLevel(spark._jvm.Level.WARN) print( "Starting collect.py with args %s" % (sys.argv) ) print("COMPACTION:input_directory:%s:output_directory:%s" % (input_directory,output_directory)) df = spark.read.csv(input_directory) cnt_before_compaction = df.count() print("INPUT DATAFRAME:rowcount:%d" % cnt_before_compaction) df.repartition(1).write.mode("overwrite").format("csv").save(output_directory) df1 = spark.read.csv(output_directory) cnt_after_compaction = df1.count() print("OUTPUT DATAFRAME:rowcount:%d" % (cnt_after_compaction)) print("End Success")
方法2:通过spark-submit参数指定日志配置
提交作业时,通过--conf参数直接设置log4j属性,无需修改代码:
spark-submit --conf "spark.driver.extraJavaOptions=-Dlog4j.logger.org.apache.spark.sql.execution.datasources.FileScanRDD=WARN" \ --conf "spark.executor.extraJavaOptions=-Dlog4j.logger.org.apache.spark.sql.execution.datasources.FileScanRDD=WARN" \ your_script.py input_dir output_dir
方法3:修改全局log4j配置文件
如果需要全局生效,找到Spark安装目录下的conf/log4j.properties(无则复制log4j.properties.template重命名),添加或修改以下配置:
log4j.logger.org.apache.spark.sql.execution.datasources.FileScanRDD=WARN
修改后重启Spark服务,后续所有作业都会应用该日志级别设置。
原理说明
这类日志是FileScanRDD类输出的INFO级消息,当读取非Hive分区格式的文件时,分区值为空就会触发该日志。通过将该类的日志级别提升至WARN及以上,就能过滤掉这些不需要的INFO消息,同时不会影响其他组件的正常日志输出。
内容的提问来源于stack exchange,提问作者Mark Rostron
相关产品推荐
相关产品推荐

