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

SparkKubernetesSensor因SparkKubernetesOperator终止触发404错误求解决方案

数据工程任务编排问题:Airflow+SparkKubernetesOperator传感器404错误

问题背景

我们的数据工程项目基于Kubernetes集群,采用SparkKubernetesOperator提交Spark任务,SparkKubernetesSensor监听任务状态,整体由Airflow编排工作流。

当前核心问题:任务竞争导致基础设施资源不足,传感器任务排队延迟过长——SparkKubernetesOperator已完成任务并终止对应Pod,而传感器启动时找不到对应的SparkApplication自定义资源,触发404错误(日志显示找不到sparkapplications.sparkoperator.k8s.io对象)。因短期内无法扩容基础设施,现寻求中短期解决方案。

已初步考虑两个方向:

  • 增加Pod的timetoliveseconds参数,提升传感器找到资源的概率,但担心闲置Pod占用节点资源引发新瓶颈。
  • 尝试让传感器与Operator同步触发,而非先后执行,但不确定实现可行性。

错误日志

[2023-05-05, 13:24:15 UTC] {taskinstance.py:1300} INFO - Executing <Task(SparkKubernetesSensor): sensor_perfil_verificado_promo> on 2023-05-04 06:00:00+00:00
[2023-05-05, 13:24:15 UTC] {standard_task_runner.py:55} INFO - Started process 20 to run task
[2023-05-05, 13:24:15 UTC] {standard_task_runner.py:82} INFO - Running: ['airflow', 'tasks', 'run', 'pipeline_best_choices', 'sensor_perfil_verificado_promo', 'scheduled__2023-05-04T06:00:00+00:00', '--job-id', '1280973', '--raw', '--subdir', 'DAGS_FOLDER/best-choices/dag-best-choices.py', '--cfg-path', '/tmp/tmpqeronp7b']
[2023-05-05, 13:24:15 UTC] {standard_task_runner.py:83} INFO - Job 1280973: Subtask sensor_perfil_verificado_promo
[2023-05-05, 13:24:15 UTC] {task_command.py:388} INFO - Running <TaskInstance: pipeline_best_choices.sensor_perfil_verificado_promo scheduled__2023-05-04T06:00:00+00:00 [running]> on host pipeline-best-choices-sensor-p-3534f4895f7a478a812050aa582b6258
[2023-05-05, 13:24:16 UTC] {logging_mixin.py:137} WARNING - /home/airflow/.local/lib/python3.7/site-packages/airflow/kubernetes/kube_config.py:39 DeprecationWarning: The delete_worker_pods option in [kubernetes] has been moved to the delete_worker_pods option in [kubernetes_executor] - the old setting has been used, but please update your config.
[2023-05-05, 13:24:16 UTC] {logging_mixin.py:137} WARNING - /home/airflow/.local/lib/python3.7/site-packages/airflow/kubernetes/kube_config.py:41 DeprecationWarning: The delete_worker_pods_on_failure option in [kubernetes] has been moved to the delete_worker_pods_on_failure option in [kubernetes_executor] - the old setting has been used, but please update your config.
[2023-05-05, 13:24:16 UTC] {logging_mixin.py:137} WARNING - /home/airflow/.local/lib/python3.7/site-packages/airflow/kubernetes/kube_config.py:44 DeprecationWarning: The worker_pods_creation_batch_size option in [kubernetes] has been moved to the worker_pods_creation_batch_size option in [kubernetes_executor] - the old setting has been used, but please update your config.
[2023-05-05, 13:24:16 UTC] {pod_generator.py:424} WARNING - Model file  does not exist
[2023-05-05, 13:24:16 UTC] {taskinstance.py:1509} INFO - Exporting the following env vars:
AIRFLOW_CTX_DAG_OWNER=
AIRFLOW_CTX_DAG_ID=pipeline_best_choices
AIRFLOW_CTX_TASK_ID=sensor_perfil_verificado_promo
AIRFLOW_CTX_EXECUTION_DATE=2023-05-04T06:00:00+00:00
AIRFLOW_CTX_TRY_NUMBER=2
AIRFLOW_CTX_DAG_RUN_ID=scheduled__2023-05-04T06:00:00+00:00
[2023-05-05, 13:24:16 UTC] {spark_kubernetes.py:105} INFO - Poking: dts-best-choices-perfil-verificado-promo-2023-05-05-12-42-06
[2023-05-05, 13:24:16 UTC] {base.py:73} INFO - Using connection ID 'kubernetes_default' for task execution.
[2023-05-05, 13:24:16 UTC] {taskinstance.py:1768} ERROR - Task failed with exception
Traceback (most recent call last):
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/providers/cncf/kubernetes/hooks/kubernetes.py", line 316, in get_custom_object
    group=group, version=version, namespace=namespace, plural=plural, name=name
  File "/home/airflow/.local/lib/python3.7/site-packages/kubernetes/client/api/custom_objects_api.py", line 1484, in get_namespaced_custom_object
    return self.get_namespaced_custom_object_with_http_info(group, version, namespace, plural, name, **kwargs)  # noqa: E501
  File "/home/airflow/.local/lib/python3.7/site-packages/kubernetes/client/api/custom_objects_api.py", line 1605, in get_namespaced_custom_object_with_http_info
    collection_formats=collection_formats)
  File "/home/airflow/.local/lib/python3.7/site-packages/kubernetes/client/api_client.py", line 353, in call_api
    _preload_content, _request_timeout, _host)
  File "/home/airflow/.local/lib/python3.7/site-packages/kubernetes/client/api_client.py", line 184, in __call_api
    _request_timeout=_request_timeout)
  File "/home/airflow/.local/lib/python3.7/site-packages/kubernetes/client/api_client.py", line 377, in request
    headers=headers)
  File "/home/airflow/.local/lib/python3.7/site-packages/kubernetes/client/rest.py", line 244, in GET
    query_params=query_params)
  File "/home/airflow/.local/lib/python3.7/site-packages/kubernetes/client/rest.py", line 234, in request
    raise ApiException(http_resp=r)
kubernetes.client.exceptions.ApiException: (404)
Reason: Not Found
HTTP response headers: HTTPHeaderDict({'Audit-Id': '41bcdf27-da7e-453b-a344-0fa82807af51', 'Cache-Control': 'no-cache, private', 'Content-Type': 'application/json', 'X-Kubernetes-Pf-Flowschema-Uid': 'ef02cbb8-3a78-4692-86aa-30128c75e9f8', 'X-Kubernetes-Pf-Prioritylevel-Uid': 'c8942bd1-20bc-4c28-b1b9-400500c532fc', 'Date': 'Fri, 05 May 2023 13:24:16 GMT', 'Content-Length': '366'})
HTTP response body: {"kind":"Status","apiVersion":"v1","metadata":{},"status":"Failure","message":"sparkapplications.sparkoperator.k8s.io \"dts-best-choices-perfil-verificado-promo-2023-05-05-12-42-06\" not found","reason":"NotFound","details":{"name":"dts-best-choices-perfil-verificado-promo-2023-05-05-12-42-06","group":"sparkoperator.k8s.io","kind":"sparkapplications"},"code":404}

During handling of the above exception, another exception occurred:

Traceback (most recent call last):
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/sensors/base.py", line 199, in execute
    poke_return = self.poke(context)
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/providers/cncf/kubernetes/sensors/spark_kubernetes.py", line 111, in poke
    namespace=self.namespace,
  File "/home/airflow/.local/lib/python3.7/site-packages/airflow/providers/cncf/kubernetes/hooks/kubernetes.py", line 320, in get_custom_object
    raise AirflowException(f"Exception when calling -> get_custom_object: {e}\n")
airflow.exceptions.AirflowException: Exception when calling -> get_custom_object: (404)
Reason: Not Found
HTTP response headers: HTTPHeaderDict({'Audit-Id': '41bcdf27-da7e-453b-a344-0fa82807af51', 'Cache-Control': 'no-cache, private', 'Content-Type': 'application/json', 'X-Kubernetes-Pf-Flowschema-Uid': 'ef02cbb8-3a78-4692-86aa-30128c75e9f8', 'X-Kubernetes-Pf-Prioritylevel-Uid': 'c8942bd1-20bc-4c28-b1b9-400500c532fc', 'Date': 'Fri, 05 May 2023 13:24:16 GMT', 'Content-Length': '366'})
HTTP response body: {"kind":"Status","apiVersion":"v1","metadata":{},"status":"Failure","message":"sparkapplications.sparkoperator.k8s.io \"dts-best-choices-perfil-verificado-promo-2023-05-05-12-42-06\" not found","reason":"NotFound","details":{"name":"dts-best-choices-perfil-verificado-promo-2023-05-05-12-42-06","group":"sparkoperator.k8s.io","kind":"sparkapplications"},"code":404}

[2023-05-05, 13:24:16 UTC] {taskinstance.py:1323} INFO - Marking task as FAILED. dag_id=pipeline_best_choices, task_id=sensor_perfil_verificado_promo, execution_date=20230504T060000, start_date=20230505T132415, end_date=20230505T132416
[2023-05-05, 13:24:16 UTC] {standard_task_runner.py:105} ERROR - Failed to execute job 1280973 for task sensor_perfil_verificado_promo (Exception when calling -> get_custom_object: (404)
Reason: Not Found
HTTP response headers: HTTPHeaderDict({'Audit-Id': '41bcdf27-da7e-453b-a344-0fa82807af51', 'Cache-Control': 'no-cache, private', 'Content-Type': 'application/json', 'X-Kubernetes-Pf-Flowschema-Uid': 'ef02cbb8-3a78-4692-86aa-30128c75e9f8', 'X-Kubernetes-Pf-Prioritylevel-Uid': 'c8942bd1-20bc-4c28-b1b9-400500c532fc', 'Date': 'Fri, 05 May 2023 13:24:16 GMT', 'Content-Length': '366'})
HTTP response body: {"kind":"Status","apiVersion":"v1","metadata":{},"status":"Failure","message":"sparkapplications.sparkoperator.k8s.io \"dts-best-choices-perfil-verificado-promo-2023-05-05-12-42-06\" not found","reason":"NotFound","details":{"name":"dts-best-choices-perfil-verificado-promo-2023-05-05-12-42-06","group":"sparkoperator.k8s.io","kind":"sparkapplications"},"code":404}

; 20)
[2023-05-05, 13:24:16 UTC] {local_task_job.py:208} INFO - Task exited with return code 1
[2023-05-05, 13:24:16 UTC] {taskinstance.py:2578} INFO - 0 downstream tasks scheduled from follow-on schedule check

中短期解决方案建议

1. 优化Pod的timetoliveseconds参数,精准控制资源占用

  • 基于Airflow任务历史统计传感器的最大排队延迟,设置TTL为该值+30-60秒,既保证传感器启动时能找到Pod/资源,又避免过度占用资源。
  • 给Spark Pod配置严格的requests和limits资源配额,限制闲置Pod的资源消耗。
  • 按任务优先级差异化设置TTL:核心业务任务TTL稍长,非核心任务保持默认值,平衡可靠性与资源利用率。

2. 调整Airflow任务依赖,让传感器提前启动

  • 无需完全同步触发,可修改依赖逻辑:将传感器任务的触发条件从Operator任务成功结束改为启动成功。具体实现:
    • 在Operator任务中添加do_xcom_push=True,将创建的SparkApplication名称推送到XCom。
    • 传感器任务依赖Operator任务的running状态(而非success),并从XCom获取SparkApplication名称进行监听。
  • 该方式能让传感器在Operator提交任务后立即启动,减少等待窗口,降低错过资源的概率。

3. 自定义容错传感器,处理404场景

  • 继承原SparkKubernetesSensor,修改poke方法,捕获404错误后进入重试逻辑,而非直接失败:
from airflow.providers.cncf.kubernetes.sensors.spark_kubernetes import SparkKubernetesSensor
from airflow.exceptions import AirflowException
from datetime import timedelta

class TolerantSparkKubernetesSensor(SparkKubernetesSensor):
    def __init__(self, *args, **kwargs):
        super().__init__(*args, **kwargs)
        # 自定义重试次数和间隔
        self.retries = 3
        self.retry_delay = timedelta(minutes=1)

    def poke(self, context):
        try:
            return super().poke(context)
        except AirflowException as e:
            if "404" in str(e) or "NotFound" in str(e):
                # 未达到重试上限则返回False,触发下一次poke
                if context['ti'].try_number <= self.retries:
                    return False
                # 重试耗尽后再抛出错误
                raise e
            else:
                raise e
  • 在DAG中使用这个自定义传感器,配合Airflow的重试机制,给传感器足够的时间等待资源出现。

4. 保留SparkApplication资源而非Pod

  • 给Spark任务添加注解sparkoperator.k8s.io/keep-application-resource: "true",让Spark Operator在Pod终止后保留SparkApplication自定义资源一段时间(默认通常为1小时)。
  • 该方式仅保留K8s中的元数据资源,占用远小于保留Pod,能有效解决传感器找不到资源的404问题,无需修改Pod的TTL。

内容的提问来源于stack exchange,提问作者Michel Senra

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 21:32:06