尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

深入解析 Apache Airflow KubernetesExecutor:Pod 模板、按任务覆盖与故障恢复机制

深入解析 Apache Airflow KubernetesExecutor:Pod 模板、按任务覆盖与故障恢复机制 深入解析 Apache Airflow KubernetesExecutorPod 模板、按任务覆盖与故障恢复机制【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowApache 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 PodPod 执行任务、上报结果后终止。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 DBPod 以Succeeded状态结束Kubernetes 将记录写入 etcdAirflow Scheduler 通过 k8s watcher 线程读取到Succeeded状态。在一个由五个节点组成的分布式 Kubernetes 集群上部署 Airflow 时Worker Pod 需要能够访问 DAG 文件以执行其中的任务并需要访问元数据库同时Kubernetes Executor 专属的配置如 Worker 命名空间、镜像信息必须在 Airflow Configuration 文件中指定。此外Executor 还支持通过executor_config以按任务per-task粒度指定额外特性。从源码结构看这条链路的实现集中在 KubernetesExecutor 类中execute_async()把任务封装为KubernetesJob放入task_queueL319-L373sync()每轮从result_queue中取出 Watcher 线程产生的KubernetesResults调用_change_state()更新任务状态然后按批次从task_queue取出任务创建 PodL401-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.yamlapiVersion: 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 都可从该卷读取同一份 DAGapiVersion: 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 拉取 DAGgit-sync对应源文件 git_sync_template.yaml。它通过initContainers中的git-sync容器在 base 容器启动之前执行一次git pullapiVersion: 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]覆盖namespacedefaultWorker Pod 所在命名空间任务pod_override中的metadata.namespace优先生效in_clusterTrue是否使用集群内配置KubeConfig 逻辑见 kube_config.pydelete_worker_podsTrue任务完成后是否删除 Worker PodFalse时改为给 Pod 打 done 注解保留现场delete_worker_pods_on_failureFalse任务失败时是否同样删除 Podworker_pods_creation_batch_size1每个调度循环最多创建多少个 Podmulti_namespace_modeFalse是否跨命名空间查询 Podworker_container_repository/worker_container_tag空Worker 基础镜像组合为kube_imagerunning_pod_log_lines100从 Pod 读取运行时日志时 tail 的行数见下文日志部分使用 pod_override 按任务覆盖 Pod使用 KubernetesExecutor 时Airflow 支持在任务粒度覆盖系统默认值。方法是构造一个 KubernetesV1Pod对象并填入期望的覆盖字段作为executor_config的一部分传入任务。官方示例 DAG example_kubernetes_executor.py 中给出了两种典型用法。覆盖 base 容器的字段挂载卷要覆盖 Executor 所启动 Pod 的 base 容器创建一个只含单个容器名为base的 V1Pod然后覆盖相应字段。下面是示例 DAG 中挂载 hostPath 卷的完整写法L62-L98executor_config_volume_mount { pod_override: k8s.V1Pod( speck8s.V1PodSpec( containers[ k8s.V1Container( namebase, volume_mounts[ k8s.V1VolumeMount(mount_path/foo/, nameexample-kubernetes-test-volume) ], ) ], volumes[ k8s.V1Volume( nameexample-kubernetes-test-volume, host_pathk8s.V1HostPathVolumeSource(path/tmp/), ) ], ) ), } task(executor_configexecutor_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(fError 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-L139executor_config_sidecar { pod_override: k8s.V1Pod( speck8s.V1PodSpec( containers[ k8s.V1Container( namebase, volume_mounts[k8s.V1VolumeMount(mount_path/shared/, nameshared-empty-dir)], ), k8s.V1Container( namesidecar, imageubuntu, args[echo retrieved from mount /shared/test.txt], command[bash, -cx], volume_mounts[k8s.V1VolumeMount(mount_path/shared/, nameshared-empty-dir)], ), ], volumes[ k8s.V1Volume(nameshared-empty-dir, empty_dirk8s.V1EmptyDirVolumeSource()), ], ) ), } task(executor_configexecutor_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(fError 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_idexample_pod_template_file, scheduleNone, start_datependulum.datetime(2021, 1, 1, tzUTC), catchupFalse, tags[example3], ) as dag: executor_config_template { pod_template_file: os.path.join(AIRFLOW_HOME, pod_templates/basic_template.yaml), pod_override: k8s.V1Pod(metadatak8s.V1ObjectMeta(labels{release: stable})), } task(executor_configexecutor_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(containerbase, tail_linesRUNNING_POD_LOG_LINES)直接通过 kube API 读取 base 容器末尾若干行日志——这也再次印证了为什么模板中 Worker 容器必须命名为base。与 CeleryExecutor 的对比与 CeleryExecutor 相比KubernetesExecutor 不需要 Redis 这类附加组件但要求能访问 Kubernetes 集群同时 Pod 的监控可以直接使用 Kubernetes 自带的监控能力。资源利用率KubernetesExecutor 中每个任务独占一个 PodPod 在任务入队时创建、任务完成时销毁。历史上在弹性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 中。此外无论使用什么 ExecutorKubernetesPodOperator 都能达到类似把任务跑进 Kubernetes的效果。故障容错与恢复机制处理分布式系统时必须假定任何组件都可能在任何时刻崩溃原因从 OOM 到节点升级不一而足。KubernetesExecutor 的容错设计围绕两个故障点展开Worker Pod 崩溃和 Scheduler 崩溃。调试利器generate-dag-yaml官方文档给出的排障建议是遇到 KubernetesExecutor 问题时可以使用airflow kubernetes generate-dag-yaml命令。该命令会按实际启动时的逻辑生成每个任务的 Pod 定义并 dump 成 YAML 文件供你检视——无需真正在集群里拉起 Pod。从 CLI 定义 可以看到它接受的参数--dag-id、--logical-date2.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 Podcleanup_pods 实现参数包括--namespace命名空间默认取[kubernetes_executor] namespace--min-pending-minutesPending 超过多少分钟的 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 集群时每条事件都带有一个单调递增的resourceVersionExecutor 每读到一个resourceVersion都会把最新值存入后端数据库因此 Scheduler 重启后可以从上次保存的resourceVersion处继续读取 watcher 事件流不会漏掉事件。从源码看这一断点续读落在 sync() 的末尾每轮把各命名空间最新的resource_version写入ResourceVersion模型元数据库表。而Scheduler 故障不会导致任务失败或重跑的根本原因在于任务独立于 Executor 运行Worker 直接把结果报告给数据库Scheduler 只是状态的观察者与更新者。此外源码还体现了更细粒度的自愈设计属于从源码结构可推断的工程细节Pod 收养adoptiontry_adopt_task_instances()/_adopt_completed_pods()会查找仍存活但原 Scheduler 已死的 Pod通过 patchairflow-workerlabel 把它们收归当前 Scheduler 的 watcher 继续管理L926-L1200避免 Scheduler 重启后出现无人跟踪的孤儿 PodPod 启动失败重试_change_state()中对任务进程尚未启动 Pod 就失败TI 仍为 QUEUED 且容器原因不在排除列表内的情况会在不消耗任务级 retry 的前提下把任务重新入队重试次数由pod_launch_failure_retries控制默认 1。小结KubernetesExecutor 的核心设计可以用三句话概括每任务一 Pod 的隔离执行模型、pod_template_filepod_override两级 Pod 定制能力、以及基于 watcher 线程和resourceVersion的故障自愈机制。落地时建议按以下顺序检查确认元数据库为非 SQLiteScheduler 可访问目标集群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),仅供参考
返回列表