Airflow Dataproc算子报错排查:任务提交与批处理创建问题
Airflow Dataproc算子两类报错解决方案
问题1:DataprocSubmitJobOperator返回"Cluster name is required"
报错信息
google.api_core.exceptions.InvalidArgument: 400 Cluster name is required
完整错误日志:
[2025-03-04, 21:37:45 UTC] {taskinstance.py:1938} ERROR - Task failed with exception Traceback (most recent call last): File "/usr/local/lib/python3.11/site-packages/google/api_core/grpc_helpers.py", line 75, in error_remapped_callable return callable_(*args, **kwargs) ^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/usr/local/lib/python3.11/site-packages/grpc/_channel.py", line 1161, in __call__ return _end_unary_response_blocking(state, call, False, None) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/usr/local/lib/python3.11/site-packages/grpc/_channel.py", line 1004, in _end_unary_response_blocking raise _InactiveRpcError(state) # pytype: disable=not-instantiable ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ grpc._channel._InactiveRpcError: <_InactiveRpcError of RPC that terminated with: status = StatusCode.INVALID_ARGUMENT details = "Cluster name is required" debug_error_string = "UNKNOWN:Error received from peer ipv4:199.36.153.8:443 {created_time:"2025-03-04T21:37:45.543592704+00:00", grpc_status:3, grpc_message:"Cluster name is required"}" > The above exception was the direct cause of the following exception: Traceback (most recent call last): File "/usr/local/lib/python3.11/site-packages/airflow/providers/google/cloud/operators/dataproc.py", line 2211, in execute job_object = self.hook.submit_job( ^^^^^^^^^^^^^^^^^^^^^ File "/usr/local/lib/python3.11/site-packages/airflow/providers/google/common/hooks/base_google.py", line 475, in inner_wrapper return func(self, *args, **kwargs) ^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/usr/local/lib/python3.11/site-packages/airflow/providers/google/cloud/hooks/dataproc.py", line 790, in submit_job return client.submit_job( ^^^^^^^^^^^^^^^^^^ File "/usr/local/lib/python3.11/site-packages/google/cloud/dataproc_v1/services/job_controller/client.py", line 547, in submit_job response = rpc( ^^^^ File "/usr/local/lib/python3.11/site-packages/google/api_core/gapic_v1/method.py", line 131, in __call__ return wrapped_func(*args, **kwargs) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/usr/local/lib/python3.11/site-packages/google/api_core/retry.py", line 366, in retry_wrapped_func return retry_target( ^^^^^^^^^^^^^ File "/usr/local/lib/python3.11/site-packages/google/api_core/retry.py", line 204, in retry_target return target() ^^^^^^^^ File "/usr/local/lib/python3.11/site-packages/google/api_core/grpc_helpers.py", line 77, in error_remapped_callable raise exceptions.from_grpc_error(exc) from exc google.api_core.exceptions.InvalidArgument: 400 Cluster name is required
问题代码
from airflow.providers.google.cloud.operators.dataproc import DataprocSubmitJobOperator dp_submit_job = DataprocSubmitJobOperator( task_id="dp_submit_job", gcp_conn_id='', #conn id for af instace project_id='', #has gcp project ID region='us-east4', job={'spark_job': {'main_class': 'WordCount', 'jar_file_uris': 'gs://<gcs-bucket-name>/wordcount_2.12-0.1.0-SNAPSHOT.jar'}}, #jobPlacement={'clusterName' : 'cluster-name'}, request_id='48c07706-feb6-4aaa-9df0-c95ccc1e2b99', #to ignore duplicate submit job requests )
解决方案
Dataproc提交任务必须指定运行集群,需将placement字段添加到job字典中(而非单独作为算子参数),且字段名需使用GCP API标准的下划线格式cluster_name:
from airflow.providers.google.cloud.operators.dataproc import DataprocSubmitJobOperator dp_submit_job = DataprocSubmitJobOperator( task_id="dp_submit_job", gcp_conn_id='', #conn id for af instace project_id='', #has gcp project ID region='us-east4', job={ 'spark_job': { 'main_class': 'WordCount', 'jar_file_uris': 'gs://<gcs-bucket-name>/wordcount_2.12-0.1.0-SNAPSHOT.jar' }, 'placement': { 'cluster_name': 'cluster-name' } }, request_id='48c07706-feb6-4aaa-9df0-c95ccc1e2b99', #to ignore duplicate submit job requests )
问题2:DataprocCreateBatchOperator返回"networks/default not found"
报错信息
google.api_core.exceptions.NotFound: 404 The resource 'projects/<project-name>/global/networks/default' was not found
完整错误日志:
Traceback (most recent call last): File "/usr/local/lib/python3.11/site-packages/google/api_core/grpc_helpers.py", line 75, in error_remapped_callable return callable_(*args, **kwargs) ^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/usr/local/lib/python3.11/site-packages/grpc/_channel.py", line 1161, in __call__ return _end_unary_response_blocking(state, call, False, None) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/usr/local/lib/python3.11/site-packages/grpc/_channel.py", line 1004, in _end_unary_response_blocking raise _InactiveRpcError(state) # pytype: disable=not-instantiable ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ grpc._channel._InactiveRpcError: <_InactiveRpcError of RPC that terminated with: status = StatusCode.NOT_FOUND details = "The resource 'projects/<project-name>/global/networks/default' was not found" debug_error_string = "UNKNOWN:Error received from peer ipv4:199.36.153.8:443 {grpc_message:"The resource \'projects/<project-name>/global/networks/default\' was not found", grpc_status:5, created_time:"2025-03-10T19:18:21.022463126+00:00"}" > The above exception was the direct cause of the following exception: Traceback (most recent call last): File "/usr/local/lib/python3.11/site-packages/airflow/providers/google/cloud/operators/dataproc.py", line 2522, in execute self.operation = hook.create_batch( ^^^^^^^^^^^^^^^^^^ File "/usr/local/lib/python3.11/site-packages/airflow/providers/google/common/hooks/base_google.py", line 475, in inner_wrapper return func(self, *args, **kwargs) ^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/usr/local/lib/python3.11/site-packages/airflow/providers/google/cloud/hooks/dataproc.py", line 863, in create_batch result = client.create_batch( ^^^^^^^^^^^^^^^^^^^^ File "/usr/local/lib/python3.11/site-packages/google/cloud/dataproc_v1/services/batch_controller/client.py", line 601, in create_batch response = rpc( ^^^^ File "/usr/local/lib/python3.11/site-packages/google/api_core/gapic_v1/method.py", line 131, in __call__ return wrapped_func(*args, **kwargs) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/usr/local/lib/python3.11/site-packages/google/api_core/grpc_helpers.py", line 77, in error_remapped_callable raise exceptions.from_grpc_error(exc) from exc google.api_core.exceptions.NotFound: 404 The resource 'projects/<project-name>/global/networks/default' was not found
问题代码
from airflow.providers.google.cloud.operators.dataproc import DataprocCreateBatchOperator # type: ignore # noqa: I001 dp_create_batch = DataprocCreateBatchOperator( task_id="dp_create_batch", gcp_conn_id='<conn-id>', project_id='<project-id>', region='us-east4', batch={ "spark_batch": { 'main_class': 'WordCount', 'jar_file_uris': ['gs://<bucket-name>/wordcount_2.12-0.1.0-SNAPSHOT.jar'], }, }, batch_id='testingbatchoperator', )
解决方案
错误原因是Dataproc Batch默认尝试使用项目的defaultVPC网络,但该网络不存在,可通过以下两种方式解决:
- 指定已存在的自定义网络
在batch字典中添加network_uri字段,指向项目中已有的VPC网络路径:
from airflow.providers.google.cloud.operators.dataproc import DataprocCreateBatchOperator # type: ignore # noqa: I001 dp_create_batch = DataprocCreateBatchOperator( task_id="dp_create_batch", gcp_conn_id='<conn-id>', project_id='<project-id>', region='us-east4', batch={ "spark_batch": { 'main_class': 'WordCount', 'jar_file_uris': ['gs://<bucket-name>/wordcount_2.12-0.1.0-SNAPSHOT.jar'], }, "network_uri": "projects/<project-id>/global/networks/<your-network-name>" }, batch_id='testingbatchoperator', )
- 创建默认VPC网络
若无需自定义网络,可在GCP控制台的VPC网络页面点击"创建默认网络",系统会自动生成default网络及配套的子网、防火墙规则。
内容的提问来源于stack exchange,提问作者Abhijit Aravind
相关产品推荐
相关产品推荐

