如何从PySpark/Python屏蔽Scala/Spark的Parquet文件JVM端异常
如何抑制Spark读取Parquet时的JVM端异常输出
问题背景
我已经通过以下代码屏蔽了py4j端的Parquet文件读取异常:
logging.getLogger('py4j').setLevel(CRITICAL) spark.read.parquet(in_path) except Exception as e: logging.getLogger('py4j').setLevel(WARN)
但JVM端仍然会打印初始异常信息:
org.apache.spark.SparkException: Exception thrown in awaitResult: at org.apache.spark.util.ThreadUtils$.awaitResult(ThreadUtils.scala:301) at org.apache.spark.util.ThreadUtils$.parmap(ThreadUtils.scala:375) at org.apache.spark.sql.execution.datasources.parquet.ParquetFileFormat$.readParquetFootersInParallel(ParquetFileFormat.scala:476) at org.apache.spark.sql.execution.datasources.parquet.ParquetFileFormat$.$anonfun$mergeSchemasInParallel$1(ParquetFileFormat.scala:522) at org.apache.spark.sql.execution.datasources.parquet.ParquetFileFormat$.$anonfun$mergeSchemasInParallel$1$adapted(ParquetFileFormat.scala:516) at org.apache.spark.sql.execution.datasources.SchemaMergeUtils$.$anonfun$mergeSchemasInParallel$2(SchemaMergeUtils.scala:76) at org.apache.spark.rdd.RDD.$anonfun$mapPartitions$2(RDD.scala:855) at org.apache.spark.rdd.RDD.$anonfun$mapPartitions$2$adapted(RDD.scala:855) at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52) at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:365) at org.apache.spark.rdd.RDD.iterator(RDD.scala:329) at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90) at org.apache.spark.scheduler.Task.run(Task.scala:136) at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:548) at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1504) at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:551) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628) at java.base/java.lang.Thread.run(Thread.java:829) Caused by: java.io.IOException: Could not read footer for file:
需要同时抑制JVM端的异常输出。
解决方案
方法1:初始化SparkSession时配置JVM日志级别
直接在SparkSession构建时指定Parquet相关JVM日志的级别,从根源限制输出:
from pyspark.sql import SparkSession # 针对log4j的配置 spark = SparkSession.builder \ .appName("SuppressParquetJVMLogs") \ .config("spark.driver.extraJavaOptions", "-Dlog4j.logger.org.apache.spark.sql.execution.datasources.parquet=ERROR") \ .config("spark.executor.extraJavaOptions", "-Dlog4j.logger.org.apache.spark.sql.execution.datasources.parquet=ERROR") \ .getOrCreate() # 如果使用logback,替换为以下配置 # spark = SparkSession.builder \ # .appName("SuppressParquetJVMLogs") \ # .config("spark.driver.extraJavaOptions", "-Dlogback.logger.org.apache.spark.sql.execution.datasources.parquet=ERROR") \ # .config("spark.executor.extraJavaOptions", "-Dlogback.logger.org.apache.spark.sql.execution.datasources.parquet=ERROR") \ # .getOrCreate()
该配置会让JVM仅输出Parquet模块的ERROR及以上级别日志,你遇到的这类IOException会被屏蔽。
方法2:动态调整JVM日志级别(无需重启Spark)
如果已经创建了SparkSession,可以通过调用JVM的日志API临时修改级别:
import logging # 获取JVM端的LogManager和目标Logger log_manager = spark._jvm.org.apache.log4j.LogManager parquet_logger = log_manager.getLogger("org.apache.spark.sql.execution.datasources.parquet") # 设置为ERROR级别 parquet_logger.setLevel(spark._jvm.org.apache.log4j.Level.ERROR) # 执行Parquet读取操作 try: logging.getLogger('py4j').setLevel(logging.CRITICAL) df = spark.read.parquet(in_path) except Exception as e: # 可选:恢复日志级别 logging.getLogger('py4j').setLevel(logging.WARN) parquet_logger.setLevel(spark._jvm.org.apache.log4j.Level.WARN)
适合临时调整日志级别的场景,无需重新初始化SparkSession。
方法3:修改全局日志配置文件(集群环境)
在Spark集群的配置文件中修改日志级别,全局生效:
- log4j配置:找到Spark安装目录下的
conf/log4j.properties(若不存在则复制log4j.properties.template并重命名),添加一行:
log4j.logger.org.apache.spark.sql.execution.datasources.parquet=ERROR
- logback配置:修改
conf/logback.xml,添加对应的logger节点:
<logger name="org.apache.spark.sql.execution.datasources.parquet" level="ERROR"/>
修改后重启Spark集群即可生效。
注意事项
- 建议精准定位到
org.apache.spark.sql.execution.datasources.parquet包调整日志级别,避免屏蔽其他Spark模块的关键日志。 - 可以通过
spark.conf.get("spark.logConf")查看当前Spark使用的日志框架类型(log4j或logback)。
内容的提问来源于stack exchange,提问作者WestCoastProjects
相关产品推荐
相关产品推荐

