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任务中添加
- 该方式能让传感器在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
相关产品推荐
相关产品推荐

