You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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网络,但该网络不存在,可通过以下两种方式解决:

  1. 指定已存在的自定义网络
    在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',
)
  1. 创建默认VPC网络
    若无需自定义网络,可在GCP控制台的VPC网络页面点击"创建默认网络",系统会自动生成default网络及配套的子网、防火墙规则。

内容的提问来源于stack exchange,提问作者Abhijit Aravind

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.14 03:23:11