深入解析 Apache Airflow KubernetesExecutor:Pod 模板、按任务覆盖与故障恢复机制
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
Apache Airflow 的 KubernetesExecutor 将每个任务实例(Task Instance)运行在 Kubernetes 集群中独立的 Pod 里,是大规模、资源异构工作负载的主流执行方案。本文基于 cncf.kubernetes provider 的官方文档与源码,完整讲解其执行模型、pod_template_file与pod_override两套 Pod 定制机制、DAG 与日志的分发方案,以及 Watcher 线程与resourceVersion支撑的故障容错原理,帮助你在生产环境正确配置、调试并恢复该执行器。
KubernetesExecutor 的执行模型
KubernetesExecutor 的工作方式可以概括为三点(见 kubernetes_executor.rst):
- 每个任务一个 Pod。当 DAG 提交一个任务时,Executor 向 Kubernetes API 申请一个 Worker Pod;Pod 执行任务、上报结果后终止。Pod 在任务入队时创建,任务完成即销毁,生命周期与任务一一对应。
- Executor 运行在 Scheduler 进程内。KubernetesExecutor 作为 Airflow Scheduler 进程内的一个模块运行;Scheduler 本身不一定要跑在 Kubernetes 上,但必须能访问一个 Kubernetes 集群。
- 要求非 SQLite 元数据库。KubernetesExecutor 需要后端使用非 sqlite 的数据库(如 PostgreSQL)来存储任务状态。
正常执行链路(Happy Path)
任务执行的完整闭环如下:
- Airflow Scheduler 向 Kubernetes 请求一个 Pod,携带
airflow run ...命令; - Kubernetes 创建 Airflow Worker Pod 并执行该命令;
- Worker 将任务成功或失败的结果写入元数据库(Airflow DB);
- Pod 以
Succeeded状态结束,Kubernetes 将记录写入 etcd; - Airflow Scheduler 通过 k8s watcher 线程读取到
Succeeded状态。
在一个由五个节点组成的分布式 Kubernetes 集群上部署 Airflow 时,Worker Pod 需要能够访问 DAG 文件以执行其中的任务,并需要访问元数据库;同时,Kubernetes Executor 专属的配置(如 Worker 命名空间、镜像信息)必须在 Airflow Configuration 文件中指定。此外,Executor 还支持通过executor_config以按任务(per-task)粒度指定额外特性。
从源码结构看,这条链路的实现集中在 KubernetesExecutor 类中:
execute_async()把任务封装为KubernetesJob放入task_queue(L319-L373);sync()每轮从result_queue中取出 Watcher 线程产生的KubernetesResults,调用_change_state()更新任务状态,然后按批次从task_queue取出任务创建 Pod(L401-L473);- 当
async_pod_creation开启时走_create_pods_concurrently()(异步客户端并发创建),否则走_create_pods_sequentially()顺序创建,每轮批量大小由worker_pods_creation_batch_size控制(L475-L540)。
start()方法中通过multiprocessing.Manager()创建进程间共享的 JoinableQueue,并启动AirflowKubernetesScheduler(即负责监视 Pod 生命周期的 Watcher 调度器)(L243-L263)。
配置 pod_template_file
要定制 KubernetesExecutor Worker Pod,可以创建 Pod 模板文件,并在airflow.cfg的[kubernetes_executor]段中通过pod_template_file选项指定其路径。provider 的默认值清单中该选项默认为空(见 provider_config_fallback_defaults.cfg)。
Airflow 对模板文件有两条硬性要求:
硬性要求一:base 容器
模板的spec.containers[0]位置必须有一个名为base的容器,且必须指定image。你可以在它之后自由添加 sidecar 容器,但 Airflow 假定 Worker 容器位于容器数组的开头且命名为base。
需要注意:Airflow 可能会覆盖 base 容器的image(例如通过下文pod_override配置),但该字段必须存在于模板文件中且不能为空。
硬性要求二:Pod 名称
模板文件中必须设置metadata.name。这个字段在每次 Pod 启动时都会被动态覆盖以保证所有 Pod 名称唯一,但它必须出现在模板中、不能留空——这正是上面三个示例模板中都写有placeholder-name/dummy-name的原因。
以下示例模板在默认 Airflow 配置下即可工作。但要注意:许多自定义配置值也需要显式地通过模板传递给 Pod,包括但不限于 SQL 连接配置、必需的 Airflow Connections、DAG 目录路径和日志设置。
示例一:将 DAG 打进镜像
对应源文件 dags_in_image_template.yaml:
apiVersion: v1 kind: Pod metadata: name: placeholder-name spec: containers: - env: - name: AIRFLOW__CORE__EXECUTOR value: LocalExecutor # Hard Coded Airflow Envs - name: AIRFLOW__CORE__FERNET_KEY valueFrom: secretKeyRef: name: RELEASE-NAME-fernet-key key: fernet-key - name: AIRFLOW__DATABASE__SQL_ALCHEMY_CONN valueFrom: secretKeyRef: name: RELEASE-NAME-airflow-metadata key: connection - name: AIRFLOW_CONN_AIRFLOW_DB valueFrom: secretKeyRef: name: RELEASE-NAME-airflow-metadata key: connection image: dummy_image imagePullPolicy: IfNotPresent name: base volumeMounts: - mountPath: "/opt/airflow/logs" name: airflow-logs - mountPath: /opt/airflow/airflow.cfg name: airflow-config readOnly: true subPath: airflow.cfg restartPolicy: Never securityContext: runAsUser: 50000 fsGroup: 50000 serviceAccountName: "RELEASE-NAME-worker-serviceaccount" volumes: - emptyDir: {} name: airflow-logs - configMap: name: RELEASE-NAME-airflow-config name: airflow-config该模板假设 DAG 已随镜像发布,因此只挂载了日志卷和airflow.cfg配置。
示例二:DAG 存放在持久卷(persistentVolume)
对应源文件 dags_in_volume_template.yaml。与示例一的差别在于多挂载了一个指向 PVC 的 DAG 卷,所有 Worker Pod 都可从该卷读取同一份 DAG:
apiVersion: v1 kind: Pod metadata: name: placeholder-name spec: containers: - env: - name: AIRFLOW__CORE__EXECUTOR value: LocalExecutor # Hard Coded Airflow Envs - name: AIRFLOW__CORE__FERNET_KEY valueFrom: secretKeyRef: name: RELEASE-NAME-fernet-key key: fernet-key - name: AIRFLOW__DATABASE__SQL_ALCHEMY_CONN valueFrom: secretKeyRef: name: RELEASE-NAME-airflow-metadata key: connection - name: AIRFLOW_CONN_AIRFLOW_DB valueFrom: secretKeyRef: name: RELEASE-NAME-airflow-metadata key: connection image: dummy_image imagePullPolicy: IfNotPresent name: base volumeMounts: - mountPath: "/opt/airflow/logs" name: airflow-logs - mountPath: /opt/airflow/dags name: airflow-dags readOnly: true - mountPath: /opt/airflow/airflow.cfg name: airflow-config readOnly: true subPath: airflow.cfg restartPolicy: Never securityContext: runAsUser: 50000 fsGroup: 50000 serviceAccountName: "RELEASE-NAME-worker-serviceaccount" volumes: - name: airflow-dags persistentVolumeClaim: claimName: RELEASE-NAME-dags - emptyDir: {} name: airflow-logs - configMap: name: RELEASE-NAME-airflow-config name: airflow-config示例三:从 git 拉取 DAG(git-sync)
对应源文件 git_sync_template.yaml。它通过initContainers中的git-sync容器在 base 容器启动之前执行一次git pull:
apiVersion: v1 kind: Pod metadata: name: dummy-name spec: initContainers: - name: git-sync image: "registry.k8s.io/git-sync/git-sync:v3.6.3" env: - name: GIT_SYNC_BRANCH value: "v2-2-stable" - name: GIT_SYNC_REPO value: "https://github.com/apache/airflow.git" - name: GIT_SYNC_DEPTH value: "1" - name: GIT_SYNC_ROOT value: "/git" - name: GIT_SYNC_DEST value: "repo" - name: GIT_SYNC_ADD_USER value: "true" - name: GIT_SYNC_ONE_TIME value: "true" volumeMounts: - name: airflow-dags mountPath: /git containers: - env: - name: AIRFLOW__CORE__EXECUTOR value: LocalExecutor # Hard Coded Airflow Envs - name: AIRFLOW__CORE__FERNET_KEY valueFrom: secretKeyRef: name: RELEASE-NAME-fernet-key key: fernet-key - name: AIRFLOW__DATABASE__SQL_ALCHEMY_CONN valueFrom: secretKeyRef: name: RELEASE-NAME-airflow-metadata key: connection - name: AIRFLOW_CONN_AIRFLOW_DB valueFrom: secretKeyRef: name: RELEASE-NAME-airflow-metadata key: connection image: dummy_image imagePullPolicy: IfNotPresent name: base volumeMounts: - mountPath: "/opt/airflow/logs" name: airflow-logs - mountPath: /opt/airflow/dags name: airflow-dags subPath: repo/airflow/example_dags readOnly: false - mountPath: /opt/airflow/airflow.cfg name: airflow-config readOnly: true subPath: airflow.cfg restartPolicy: Never securityContext: runAsUser: 50000 fsGroup: 50000 serviceAccountName: "RELEASE-NAME-worker-serviceaccount" volumes: - name: airflow-dags emptyDir: {} - name: airflow-logs emptyDir: {} - configMap: name: RELEASE-NAME-airflow-config name: airflow-config[kubernetes_executor] 关键配置项速览
结合默认配置文件与 Executor 源码中的实际读取逻辑,以下是最常用的[kubernetes_executor]配置项:
| 配置项 | 默认值 | 说明 |
|---|---|---|
pod_template_file | 空 | 全局 Pod 模板文件路径,可被任务级executor_config["pod_template_file"]覆盖 |
namespace | default | Worker Pod 所在命名空间;任务pod_override中的metadata.namespace优先生效 |
in_cluster | True | 是否使用集群内配置(KubeConfig 逻辑见 kube_config.py) |
delete_worker_pods | True | 任务完成后是否删除 Worker Pod;False时改为给 Pod 打 done 注解保留现场 |
delete_worker_pods_on_failure | False | 任务失败时是否同样删除 Pod |
worker_pods_creation_batch_size | 1 | 每个调度循环最多创建多少个 Pod |
multi_namespace_mode | False | 是否跨命名空间查询 Pod |
worker_container_repository/worker_container_tag | 空 | Worker 基础镜像,组合为kube_image |
running_pod_log_lines | 100 | 从 Pod 读取运行时日志时 tail 的行数(见下文日志部分) |
使用 pod_override 按任务覆盖 Pod
使用 KubernetesExecutor 时,Airflow 支持在任务粒度覆盖系统默认值。方法是构造一个 KubernetesV1Pod对象并填入期望的覆盖字段,作为executor_config的一部分传入任务。官方示例 DAG example_kubernetes_executor.py 中给出了两种典型用法。
覆盖 base 容器的字段(挂载卷)
要覆盖 Executor 所启动 Pod 的 base 容器,创建一个只含单个容器(名为base)的 V1Pod,然后覆盖相应字段。下面是示例 DAG 中挂载 hostPath 卷的完整写法(L62-L98):
executor_config_volume_mount = { "pod_override": k8s.V1Pod( spec=k8s.V1PodSpec( containers=[ k8s.V1Container( name="base", volume_mounts=[ k8s.V1VolumeMount(mount_path="/foo/", name="example-kubernetes-test-volume") ], ) ], volumes=[ k8s.V1Volume( name="example-kubernetes-test-volume", host_path=k8s.V1HostPathVolumeSource(path="/tmp/"), ) ], ) ), } @task(executor_config=executor_config_volume_mount) def test_volume_mount(): """ Tests whether the volume has been mounted. """ with open("/foo/volume_mount_test.txt", "w") as foo: foo.write("Hello") return_code = os.system("cat /foo/volume_mount_test.txt") if return_code != 0: raise ValueError(f"Error when checking volume mount. Return code {return_code}") volume_task = test_volume_mount()注意:以下字段会被扩展(extend)而不是覆盖(overwrite)——来自spec的volumes和init_containers;来自容器的volume_mounts、环境变量、ports和devices。
添加 sidecar 容器
要给 Pod 添加 sidecar 容器,创建一个 V1Pod:第一个容器留空但命名为base,第二个容器即为你期望的 sidecar。示例 DAG 中的完整写法(L100-L139):
executor_config_sidecar = { "pod_override": k8s.V1Pod( spec=k8s.V1PodSpec( containers=[ k8s.V1Container( name="base", volume_mounts=[k8s.V1VolumeMount(mount_path="/shared/", name="shared-empty-dir")], ), k8s.V1Container( name="sidecar", image="ubuntu", args=['echo "retrieved from mount" > /shared/test.txt'], command=["bash", "-cx"], volume_mounts=[k8s.V1VolumeMount(mount_path="/shared/", name="shared-empty-dir")], ), ], volumes=[ k8s.V1Volume(name="shared-empty-dir", empty_dir=k8s.V1EmptyDirVolumeSource()), ], ) ), } @task(executor_config=executor_config_sidecar) def test_sharedvolume_mount(): """ Tests whether the volume has been mounted. """ for i in range(5): try: return_code = os.system("cat /shared/test.txt") if return_code != 0: raise ValueError(f"Error when checking volume mount. Return code {return_code}") except ValueError as e: if i > 4: raise e sidecar_task = test_sharedvolume_mount()任务级 pod_template_file
还可以在任务粒度指定自定义的pod_template_file,以便在多个任务之间复用同一份基础值。它会替换airflow.cfg中指定的默认模板,然后再用pod_override进行覆盖;该模板同时会被用来生成 Airflow UI 中可见的 Pod K8s Spec。两者结合使用的完整示例:
import os import pendulum from airflow import DAG from airflow.decorators import task from airflow.example_dags.libs.helper import print_stuff from airflow.settings import AIRFLOW_HOME from kubernetes.client import models as k8s with DAG( dag_id="example_pod_template_file", schedule=None, start_date=pendulum.datetime(2021, 1, 1, tz="UTC"), catchup=False, tags=["example3"], ) as dag: executor_config_template = { "pod_template_file": os.path.join(AIRFLOW_HOME, "pod_templates/basic_template.yaml"), "pod_override": k8s.V1Pod(metadata=k8s.V1ObjectMeta(labels={"release": "stable"})), } @task(executor_config=executor_config_template) def task_with_template(): print_stuff()从源码结构看,任务级模板与覆盖对象在execute_async()中通过PodGenerator.from_obj(executor_config)统一解析为kube_executor_config,并连同pod_template_file路径一起放入KubernetesJob;最终 Pod 由 PodGenerator.construct_pod 基于base_worker_pod(模板反序列化结果)+pod_override_object合并构造。
管理 DAG 与日志
是否使用持久卷取决于你的配置。
DAG 的三种分发方式:
- 将 DAG 直接包含在镜像中;
- 使用
git-sync:在 base Worker 容器启动前执行一次git pull拉取 DAG 仓库; - 将 DAG 存放在持久卷上,并挂载到所有 Worker。
日志的两条出路:
- 使用持久卷,同时挂载到 Webserver 和 Worker 上;
- 启用远程日志(remote logging)。
注意:如果你既不启用日志持久化、也未启用远程日志,日志将在 Worker Pod 关闭后丢失。
源码层面,UI 对 RUNNING 状态任务读取 Pod 日志的实现在 get_streaming_task_log:它按dag_id/task_id/run_id/try_number等 label 精确定位 Pod,然后调用read_namespaced_pod_log(container="base", tail_lines=RUNNING_POD_LOG_LINES)直接通过 kube API 读取 base 容器末尾若干行日志——这也再次印证了为什么模板中 Worker 容器必须命名为base。
与 CeleryExecutor 的对比
与 CeleryExecutor 相比,KubernetesExecutor 不需要 Redis 这类附加组件,但要求能访问 Kubernetes 集群;同时 Pod 的监控可以直接使用 Kubernetes 自带的监控能力。
资源利用率:KubernetesExecutor 中每个任务独占一个 Pod,Pod 在任务入队时创建、任务完成时销毁。历史上,在弹性(burstable)负载场景下,这相对于 CeleryExecutor 具有资源利用率优势——Celery 需要你维护一组固定数量、长期运行的 Worker Pod,无论有没有任务。但需要注意:官方 Apache Airflow Helm chart 已支持根据队列中任务数量把 Celery Worker 自动缩容到零,因此在使用官方 chart 时,这已不再是一个决定性优势。
任务延迟与资源隔离:由于 Celery Worker Pod 在任务入队前就已经运行就绪,任务延迟通常更低;反之,KubernetesExecutor 需要经历拉镜像、建 Pod 的过程。另一方面,Celery 下多个任务共享同一个 Pod,任务设计时(尤其是内存消耗)必须更关注资源占用。
适用场景:
- 长时任务:KubernetesExecutor 下,如果部署发生在任务运行期间,该任务会一直运行到完成(或超时等);而 CeleryExecutor 中,任务最多只能运行到宽限期(grace period)结束,之后会被终止。
- 资源需求或镜像不统一的负载:KubernetesExecutor 可以为不同任务使用不同镜像和资源配置,天然适合这种异构场景。
两者并非互斥:使用 CeleryKubernetesExecutor 可以在同一集群上同时使用 CeleryExecutor 和 KubernetesExecutor。它根据任务的queue字段决定执行位置:默认任务发给 Celery Worker;若希望某个任务走 KubernetesExecutor,将其发往kubernetes队列即可,该任务会运行在独立 Pod 中。此外,无论使用什么 Executor,KubernetesPodOperator 都能达到类似"把任务跑进 Kubernetes"的效果。
故障容错与恢复机制
处理分布式系统时,必须假定任何组件都可能在任何时刻崩溃,原因从 OOM 到节点升级不一而足。KubernetesExecutor 的容错设计围绕两个故障点展开:Worker Pod 崩溃和 Scheduler 崩溃。
调试利器:generate-dag-yaml
官方文档给出的排障建议是:遇到 KubernetesExecutor 问题时,可以使用airflow kubernetes generate-dag-yaml命令。该命令会按实际启动时的逻辑生成每个任务的 Pod 定义,并 dump 成 YAML 文件供你检视——无需真正在集群里拉起 Pod。
从 CLI 定义 可以看到它接受的参数:--dag-id、--logical-date(2.x 版本为--execution-date)、--output-path、--team(多团队配置)、--verbose,以及 3.x 的--bundle-name/ 2.x 的--subdir。其实现 generate_pod_yaml 复用与 Executor 相同的PodGenerator.construct_pod路径(同样解析pod_template_file与executor_config中的pod_override),因此输出的 YAML 与真实启动的 Pod 高度一致。
同组下还有一个运维命令airflow kubernetes cleanup-pods,用于清理处于 evicted/failed/succeeded/pending 状态的 Airflow Pod(cleanup_pods 实现),参数包括:
--namespace:命名空间,默认取[kubernetes_executor] namespace;--min-pending-minutes:Pending 超过多少分钟的 Pod 才会被清理,默认 30,最小 5(低于 5 会被自动抬升到 5,以保护新建 Pod);--min-completed-minutes:已终态 Pod 至少保留多少分钟,默认 1,设为 0 则立即删除;--verbose:输出详细过程。
Worker Pod 崩溃:Watcher 线程兜底
当 Worker 在把状态回报给后端数据库之前就死亡时,Executor 依靠 Kubernetes watcher 线程来发现失败的 Pod。
流程是:Worker Pod 在任务完成前失败 → Pod 以Failed状态结束并被记录在 etcd → Scheduler 的 watcher 线程读到Failed→ Scheduler 在数据库中将任务记为FAILED。
所谓 Kubernetes watcher,就是一个可以订阅 Kubernetes 数据库(API Server watch 流)中所有变更的线程,Pod 的启动、运行、结束、失败都会触发事件。通过监视这条事件流,KubernetesExecutor 就能发现 Worker 崩溃并把任务正确标记为失败,而不必依赖 Worker 自身的回报。源码中该逻辑由AirflowKubernetesScheduler承载,其产出的KubernetesResults(含state、resource_version与预收集的failure_details)经result_queue交给sync()处理;_change_state()在任务失败时会把pod_reason、container_state、exit_code等细节写入告警日志(L669-L752),方便定位 OOM、镜像拉取失败等原因。
Scheduler 崩溃:resourceVersion 断点续读
如果崩溃的是 Scheduler Pod 本身,恢复依赖 watcher 的resourceVersion机制:
- 在监视 Kubernetes 集群时,每条事件都带有一个单调递增的
resourceVersion; - Executor 每读到一个
resourceVersion,都会把最新值存入后端数据库; - 因此 Scheduler 重启后,可以从上次保存的
resourceVersion处继续读取 watcher 事件流,不会漏掉事件。
从源码看,这一"断点续读"落在 sync() 的末尾:每轮把各命名空间最新的resource_version写入ResourceVersion模型(元数据库表)。而"Scheduler 故障不会导致任务失败或重跑"的根本原因在于:任务独立于 Executor 运行,Worker 直接把结果报告给数据库,Scheduler 只是状态的观察者与更新者。
此外,源码还体现了更细粒度的自愈设计(属于从源码结构可推断的工程细节):
- Pod 收养(adoption):
try_adopt_task_instances()/_adopt_completed_pods()会查找仍存活但"原 Scheduler 已死"的 Pod,通过 patchairflow-workerlabel 把它们收归当前 Scheduler 的 watcher 继续管理(L926-L1200),避免 Scheduler 重启后出现无人跟踪的孤儿 Pod; - Pod 启动失败重试:
_change_state()中对"任务进程尚未启动 Pod 就失败"(TI 仍为 QUEUED 且容器原因不在排除列表内)的情况,会在不消耗任务级 retry 的前提下把任务重新入队,重试次数由pod_launch_failure_retries控制(默认 1)。
小结
KubernetesExecutor 的核心设计可以用三句话概括:每任务一 Pod 的隔离执行模型、pod_template_file+pod_override两级 Pod 定制能力、以及基于 watcher 线程和resourceVersion的故障自愈机制。落地时建议按以下顺序检查:
- 确认元数据库为非 SQLite,Scheduler 可访问目标集群(
in_cluster或 kubeconfig); - 在
[kubernetes_executor]中配置namespace、worker_container_repository/tag与pod_template_file,模板务必满足base容器和metadata.name两条硬性要求; - 选定 DAG 分发方式(镜像 / git-sync / 持久卷)与日志出路(持久卷 / 远程日志),否则日志会在 Pod 关闭后丢失;
- 排查问题时先跑
airflow kubernetes generate-dag-yaml导出实际 Pod YAML,用airflow kubernetes cleanup-pods清理残留 Pod。
相关入口文件:KubernetesExecutor 实现、Pod 构造器、CLI 命令定义、示例 DAG。
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考