使用DataProcPySparkOperator时遇区域与可用区配置错误问题
刚看到你遇到的这个问题,我之前在整合Airflow和Dataproc的时候也碰过类似的状况,咱们一步步排查解决:
最可能的几个原因及对应解决方案
1. 集群区域与任务提交的region不匹配
你当前设置的region='global',但集群所在的zone是us-central1-b(属于us-central1区域)。Dataproc已经逐步弃用global区域的支持,且任务提交的region必须和集群创建时的区域完全一致。
修复方法:
把Operator的region参数改成集群所在的区域,修改后的代码片段:
run_pyspark_job= DataProcPySparkOperator( task_id='pyspark_operator_test', main='/root/airflow/dags/basic_eda.py', # 这里还有个容易踩坑的路径问题! job_name='test_pyspark_job', cluster_name='test-cluster-20180502', gcp_conn_id='google_cloud_default', region='us-central1', # 修改为集群所在区域 zone='us-central1-b' )
2. 集群名称错误或集群已不存在
仔细检查cluster_name的拼写是否完全正确,或者这个集群是否已经被删除。你可以通过GCP命令行工具验证集群状态:
gcloud dataproc clusters list --region=us-central1
如果列表里找不到test-cluster-20180502,要么创建对应的集群,要么把cluster_name改成现有集群的名称。
3. PySpark脚本路径错误(容易忽略的关键问题)
你当前的main参数填的是Airflow服务器的本地路径/root/airflow/dags/basic_eda.py,但Dataproc集群无法访问Airflow服务器的本地文件系统,必须把脚本上传到Google Cloud Storage(GCS),然后填写GCS路径。
修复方法:
- 把
basic_eda.py上传到你的GCS存储桶,比如gs://your-airflow-bucket/scripts/basic_eda.py - 修改Operator的main参数为GCS路径:
main='gs://your-airflow-bucket/scripts/basic_eda.py'
4. GCP连接权限不足
检查Airflow中gcp_conn_id='google_cloud_default'对应的连接配置,确保关联的服务账号拥有Dataproc Editor或至少Dataproc Worker的IAM角色,能正常调用Dataproc API并提交任务到目标集群。
验证步骤
修改完上述可能的问题后,先单独测试集群的可访问性,再运行Airflow DAG,这样能快速定位剩余问题。
内容的提问来源于stack exchange,提问作者Shrashti

