如何在Airflow的DataprocSubmitJobOperator中导入外部Python文件?
解决DataprocSubmitJobOperator导入common.zip模块的问题
问题根源
- 你错误地将
python_file_uris放在了properties字典中,这不是Spark的合法配置项,导致系统抛出Warning: Ignoring non-Spark config property: python_file_uris,zip包根本没被添加到Python路径。 archive_uris仅会将压缩包解压到作业的工作目录,但不会自动添加到Python模块搜索路径,所以单独用它无法解决导入问题。
正确配置方式
将python_file_uris从properties中移出,作为pyspark_job的顶层参数(和main_python_file_uri同级),它的作用是将指定的Python文件/zip包添加到Spark的Python路径中,让主脚本可以直接导入其中的模块。
修改后的作业配置代码:
COMMON = "gs://my-bucket/common/common.zip" INGEST_JOB_URI = "pyspark-file.py" job = { "reference": {"project_id": project_id}, "placement": {"cluster_name": cluster_name}, "pyspark_job": { "main_python_file_uri": INGEST_JOB_URI, "jar_file_uris": [DELTA_CORE_JAR_FILE_URI, GRAPHFRAME_JAR_FILE_URI, BQ_JAR_FILE_URI], "properties": { "spark.sql.extensions": "io.delta.sql.DeltaSparkSessionExtension", "spark.sql.catalog.spark_catalog": "org.apache.spark.sql.delta.catalog.DeltaCatalog" }, # 将python_file_uris移到这里,作为pyspark_job的顶层参数 "python_file_uris": [COMMON], "args": args }, }
验证说明
- 保留原导入语句
from common.file_1 import myfunction即可,因为python_file_uris会自动把common.zip加入Python模块搜索路径,Spark能识别zip包内的common目录作为合法模块。 - 如果
common.zip的结构是顶层直接包含file_1.py(而非嵌套在common目录下),则导入语句需要改为from file_1 import myfunction,请根据你的zip包实际结构调整。
内容的提问来源于stack exchange,提问作者aya62
相关产品推荐
相关产品推荐

