无法通过Python代码配置Spark日志属性,gcloud命令可正常运行
问题描述
我通过以下gcloud命令可以成功抑制Spark日志并指定自定义log4j.properties文件:
gcloud dataproc jobs submit spark \ --cluster test-dataproc-cluster \ --region europe-north1 \ --files gs://test-spark-logging-bucket/log4j.properties \ --properties spark.sql.legacy.allowUntypedScalaUDF=true,'spark.driver.extraJavaOptions=-Dlog4j.configuration=file:log4j.properties,spark.executor.extraJavaOptions=-Dlog4j.configuration=file:log4j.properties' \ --class com.pythian.edp.pm.spark.job.TestSparkJob \ --jars gs://18621ad39476e29c-test-static/spark-jobs/sparkBigQueryConnector/spark-bigquery-assembly-0.11.1-beta-SNAPSHOT.jar,gs://18621ad39476e29c-test-static/spark-jobs/digiSparkPmCmProcessing/digiSparkPmCmProcessing-assembly-0.1.0-SNAPSHOT.jar \ -- --configurationURI gs://18621ad39476e29c-test-static/spark-job-configs/f9fabc4f-162e-40f3-af69-237e4c464c9e-PM-LTE-CELLCQI/ing-15min-cell-1719203688539.yml \ --jobType "ingestion" --dataType "PM" \ --sendPubSubNotification --pubSubProjectId bmas-eu-digi-pipe-uat --pubSubTopicName data-notifications \ --traceDatasets > test_log_ingestion_job_log 2>&1
但用Python代码配置相同参数时遇到问题,以下是我的代码片段:
job_properties.update( {"spark.dynamicAllocation.enabled":"true", "spark.dynamicAllocation.minExecutors" : "0", "spark.dynamicAllocation.maxExecutors" : "5", "spark.executor.instances": "0", "spark.sql.legacy.allowUntypedScalaUDF":"true", "spark.executor.extraJavaOptions":"-Dlog4j.configuration=file:log4j.properties", "spark.sql.legacy.allowUntypedScalaUDF":"-Dlog4j.configuration=file:log4j.properties" }) # log_file_location='gs://test-spark-logging-bucket/log4j.properties' job_details = { 'placement': { 'cluster_name': cluster_name, }, 'reference': { 'job_id': job_id, }, 'scheduling': { 'max_failures_per_hour': 1, }, 'labels': labels, 'spark_job': { 'args': spark_job_arguments, 'main_class': SPARK_JOB_CLASSNAME, 'jar_file_uris': [ os.path.join(self.dataproc_job_jar_file_prefix, file) for file in jar_files ], 'properties': job_properties, } }
代码中的问题
- 重复覆盖属性:
spark.sql.legacy.allowUntypedScalaUDF被设置了两次,第二次的值-Dlog4j.configuration=file:log4j.properties会覆盖第一次的true,导致原属性失效。 - 缺失驱动端日志配置:只设置了
spark.executor.extraJavaOptions,但没有配置spark.driver.extraJavaOptions,而gcloud命令中同时配置了两者。 - 未添加log4j.properties文件:gcloud命令用
--files指定了日志文件,但Python代码中没有将该文件加入到任务的文件列表中,导致Driver和Executor找不到log4j.properties。
修正后的代码
# 修复属性重复问题,添加driver端日志配置 job_properties.update( { "spark.dynamicAllocation.enabled": "true", "spark.dynamicAllocation.minExecutors": "0", "spark.dynamicAllocation.maxExecutors": "5", "spark.executor.instances": "0", "spark.sql.legacy.allowUntypedScalaUDF": "true", "spark.driver.extraJavaOptions": "-Dlog4j.configuration=file:log4j.properties", "spark.executor.extraJavaOptions": "-Dlog4j.configuration=file:log4j.properties" } ) log_file_location = 'gs://test-spark-logging-bucket/log4j.properties' job_details = { 'placement': { 'cluster_name': cluster_name, }, 'reference': { 'job_id': job_id, }, 'scheduling': { 'max_failures_per_hour': 1, }, 'labels': labels, 'spark_job': { 'args': spark_job_arguments, 'main_class': SPARK_JOB_CLASSNAME, 'jar_file_uris': [ os.path.join(self.dataproc_job_jar_file_prefix, file) for file in jar_files ], # 添加log4j.properties到任务文件列表 'file_uris': [log_file_location], 'properties': job_properties, } }
关键说明
- file_uris字段:对应gcloud命令中的
--files参数,需要将log4j.properties的GCS路径加入,Dataproc会自动将该文件分发到Driver和Executor的工作目录,这样file:log4j.properties才能被正确找到。 - 同时配置Driver和Executor:必须为
spark.driver.extraJavaOptions和spark.executor.extraJavaOptions都设置日志配置,否则Driver端的日志不会被抑制。 - 避免属性重复:确保每个Spark属性只设置一次,防止被意外覆盖。
内容的提问来源于stack exchange,提问作者Vikrant Singh Rana
相关产品推荐
相关产品推荐

