Apache Airflow CeleryKubernetesExecutor 混合执行器全解析:原理、配置与迁移指南

发布时间:2026/9/14 13:13:51
Apache Airflow CeleryKubernetesExecutor 混合执行器全解析:原理、配置与迁移指南 Apache Airflow CeleryKubernetesExecutor 混合执行器全解析原理、配置与迁移指南【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow本文以 Apache Airflow Celery 提供者apache-airflow-providers-celery官方文档 celery_kubernetes_executor.rst 为主线深入讲解CeleryKubernetesExecutor混合执行器的设计动机、基于队列的路由机制、部署配置以及它在 Airflow 3.0 中被多执行器并发特性取代后的迁移路径。读完本文你将理解如何在同一 Airflow 环境中同时利用 Celery 的横向扩展能力和 Kubernetes 的运行时隔离并掌握在 Airflow 2.x 上配置该执行器、在 Airflow 3.x 上平滑迁移到新方案的具体做法。什么是 CeleryKubernetesExecutorCeleryKubernetesExecutor全类名airflow.providers.celery.executors.celery_kubernetes_executor.CeleryKubernetesExecutor允许用户同时运行一个CeleryExecutor和一个KubernetesExecutor并根据任务所属的**队列queue**决定由哪一个执行器来执行该任务。它的设计目标是取两者之长继承CeleryExecutor 的横向扩展能力Celery 拥有常驻 worker 集群队列消息通过 broker 分发能够在高峰期为大量任务提供稳定的执行吞吐继承KubernetesExecutor 的运行时隔离能力Kubernetes 任务以独立 Pod 运行每个任务拥有独立的容器化环境可以按任务定制镜像、依赖与资源配额避免吵闹邻居问题。从源码结构看CeleryKubernetesExecutor继承自BaseExecutor但它本质上是两个子执行器的包装器wrapper而不是一个独立的执行引擎。其核心实现在 celery_kubernetes_executor.pyclass CeleryKubernetesExecutor(BaseExecutor): CeleryKubernetesExecutor consists of CeleryExecutor and KubernetesExecutor. It chooses an executor to use based on the queue defined on the task. When the queue is the value of kubernetes_queue in section [celery_kubernetes_executor] of the configuration (default value: kubernetes), KubernetesExecutor is selected to run the task, otherwise, CeleryExecutor is used. 因此文档明确指出只有满足特定条件时才推荐使用它——因为它要求同时搭建并维护 CeleryExecutor 与 KubernetesExecutor 两套基础设施。什么时候应该使用 CeleryKubernetesExecutor原文档给出了三条明确的推荐条件三者需要结合判断高峰期的待调度任务数量超过了 Kubernetes 集群能够从容处理的规模如果全部任务都走 Kubernetes Pod 拉起流程调度高峰时可能出现 Pod 创建延迟、资源抢占等压力而 Celery 常驻 worker 可以更平稳地消化大批量任务。只有相对小部分任务需要运行时隔离运行隔离是有成本的Pod 启动有延迟、小任务跑在容器里开销偏大如果绝大多数任务并不需要隔离把它们路由给 Celery 常驻 worker 更经济。你有一大批可以在 Celery worker 上执行的小任务同时也存在资源密集、更适合跑在预定义环境中的任务典型组合是海量轻量任务走 Celery 少量重型任务走 Kubernetes 定制镜像。换言之该执行器的价值在于按队列把两类任务分流到最合适的执行环境而不是让所有任务都挤在同一种执行模型里。⚠️重要前提当前仓库Airflow 3.x 主线已经不再支持该执行器。文档开头以醒目的 note 声明从 Airflow 3.0.0 起CeleryKubernetesExecutor不再受支持官方建议改用Using Multiple Executors Concurrently多执行器并发特性后者以更灵活的方式提供等价能力。下面第 6 节会给出详细迁移说明。基于队列的路由机制核心原理这是CeleryKubernetesExecutor的灵魂用任务实例上的queue字段决定执行器。其路由逻辑集中在_router方法中celery_kubernetes_executor.pydef _router(self, simple_task_instance: SimpleTaskInstance) - CeleryExecutor | KubernetesExecutor: if simple_task_instance.queue self.kubernetes_queue: return self.kubernetes_executor return self.celery_executor即任务的queue等于kubernetes_queue时 → 交给KubernetesExecutor执行其余任何队列 → 交给CeleryExecutor执行。kubernetes_queue 配置项kubernetes_queue的值从配置段[celery_kubernetes_executor]中读取默认值为kubernetes。对应的默认配置定义在 provider_config_fallback_defaults.cfg[local_kubernetes_executor] kubernetes_queue kubernetes [celery_kubernetes_executor] kubernetes_queue kubernetes源码中的读取方式celery_kubernetes_executor.pycached_property providers_configuration_loaded def kubernetes_queue(self) - str: return conf.get(celery_kubernetes_executor, kubernetes_queue)需要注意的是构造混合执行器时KubernetesExecutor子实例会被强制同步为同一个kubernetes_queuecelery_kubernetes_executor.py保证两侧路由口径一致self.kubernetes_executor.kubernetes_queue self.kubernetes_queue单元测试 test_celery_kubernetes_executor.py 专门验证了这一行为def test_kubernetes_executor_knows_its_queue(self): ... assert k8s_executor_mock.kubernetes_queue conf.get( celery_kubernetes_executor, kubernetes_queue )队列如何被任务使用在 Airflow 中queue是BaseOperator的一个属性任何任务都可以指定自己的队列。当你想把某个任务推给 Kubernetes 侧执行时只需把它的queue设为kubernetes或你自定义的kubernetes_queue值from airflow.operators.bash import BashOperator heavy_task BashOperator( task_idheavy_job, queuekubernetes, # 走 KubernetesExecutor获得运行时隔离 bash_commandtrain_model.sh, ) light_task BashOperator( task_idlight_job, queuedefault, # 其余队列走 CeleryExecutor bash_commandecho hello, )如果不显式设置任务会落入环境默认队列配置项operators - default_queue从而默认由 CeleryExecutor 执行。配置 CeleryKubernetesExecutor由于它同时依赖两套执行器配置也分为两层1. Celery 侧配置CeleryExecutor 的完整配置参数集中在 Celery 提供者的配置参考文档中。与混合执行器直接相关的[celery]段默认值如下见 provider_config_fallback_defaults.cfg[celery] celery_app_name airflow.executors.celery_executor worker_concurrency 16 worker_prefetch_multiplier 1 worker_enable_remote_control true broker_url redis://redis:6379/0 result_backend_sqlalchemy_engine_options flower_host 0.0.0.0 flower_port 5555 flower_basic_auth sync_parallelism 0 celery_config_options airflow.providers.celery.executors.default_celery.DEFAULT_CELERY_CONFIG ssl_active False pool prefork operation_timeout 1.0 task_track_started True task_publish_max_retries 3要让 Celery 侧真正工作你还需要准备一个 Celery broker 与 result backendRabbitMQ、Redis、Redis Sentinel 等并将[celery] broker_url指向 broker按需安装依赖例如pip install apache-airflow[celery]推荐安装 Celery bundle在部署了 Celery worker 的机器上启动 workerairflow celery worker停止用airflow celery stop可选启动 Flower 监控 Web UIairflow celery flower需已安装flower库。Celery worker 也可以只监听特定队列例如airflow celery worker -q default,spark这为同一套 Celery 集群内部再分流提供了灵活性。更完整的 Celery 架构、任务执行时序与运维注意事项见 celery_executor.rst。下图展示了混合执行器中 Celery 一侧的任务执行时序——调度器将命令投递到 QueueBrokerWorkerProcess 领取任务并派发给 WorkerChildProcess后者拉起 LocalTaskJobProcess / RawTaskProcess 执行用户代码最终将状态写入 ResultBackend 供调度器轮询2. Kubernetes 侧配置Kubernetes 侧使用[kubernetes_executor]段的常规配置关键默认值见 provider_config_fallback_defaults.cfg[kubernetes_executor] pod_template_file worker_container_repository worker_container_tag namespace default delete_worker_pods True delete_worker_pods_on_failure False worker_pods_creation_batch_size 1 multi_namespace_mode False in_cluster True verify_ssl True其中pod_template_file、worker_container_repository/worker_container_tag用于定义任务 Pod 的镜像与环境这正是 Kubernetes 侧预定义环境、运行时隔离能力的来源。KubernetesExecutor 的完整说明可参考 cncf-kubernetes 提供者的 kubernetes_executor.rst。3. 启用混合执行器最后在airflow.cfg中把环境级执行器指向混合类[core] executor CeleryKubernetesExecutor从源码看Airflow 2.x 场景下构造函数接收两个子执行器实例celery_kubernetes_executor.py而在 Airflow 3.0 环境下构造函数会直接抛出RuntimeError提示改用多执行器并发特性——这一点与文档3.0 起不再支持的声明完全一致。混合执行器的运行机制与源码印证任务入队无论是queue_command还是queue_task_instance混合执行器都会先通过_router选中子执行器再转发调用celery_kubernetes_executor.py并在调试日志中记录路由结果executor self._router(task_instance) self.log.debug(Using executor: %s for %s, executor.__class__.__name__, task_instance.key)单元测试用参数化用例验证了路由的正确性test_celery_kubernetes_executor.pypytest.mark.parametrize(test_queue, [any-other-queue, KUBERNETES_QUEUE]) def test_queue_command(self, k8s_queue_cmd, celery_queue_cmd, test_queue): ... if test_queue KUBERNETES_QUEUE: k8s_queue_cmd.assert_called_once_with(simple_task_instance, *kwarg_values) celery_queue_cmd.assert_not_called() else: celery_queue_cmd.assert_called_once_with(simple_task_instance, *kwarg_values) k8s_queue_cmd.assert_not_called()生命周期与状态聚合混合执行器把生命周期方法原样转发给两个子执行器并把状态数据做并集合并start()/end()/terminate()依次调用 celery 与 kubernetes 两个子执行器的同名方法celery_kubernetes_executor.pyheartbeat()同时驱动两个子执行器的心跳celery_kubernetes_executor.pyqueued_tasks/running分别返回两个子执行器状态的并集get_event_buffer()合并并清空两个子执行器的事件缓冲job_idsetter把调度器的 job id 同步给两个子执行器测试 test_job_id_setter 验证了三者一致。任务接管与日志try_adopt_task_instances()调度器重启后尝试接管遗留任务实例会先按队列把任务拆分再分别交给对应子执行器celery_kubernetes_executor.pyget_task_log()/get_streaming_task_log()只有Kubernetes 队列上的任务才会向 KubernetesExecutor 索取日志Celery 队列任务返回空因为 Celery worker 日志通常从 worker 本地或远程日志后端获取——相关行为同样有测试覆盖test_celery_kubernetes_executor.pyrevoke_task()按队列分别调用两个子执行器的撤销逻辑并对旧版 kubernetes provider 做了cleanup_stuck_queued_tasks的回退兼容get_cli_commands()合并返回 Celery 与 Kubernetes 两套 CLI 命令celery_kubernetes_executor.py因此airflow celery worker、airflow kubernetes ...等子命令依然可用。已知限制源码可证supports_sentry False不支持 Sentry 集成serve_logs False、is_local False、is_single_threaded False、is_production True属性 setter如queued_tasks、running在混合执行器上未实现仅作为只读聚合视图。弃用说明为什么 Airflow 3.0 起不再支持Airflow 核心文档 executor/index.rst 给出了明确的弃用理由值得仔细理解实现不属于 Airflow 核心CeleryKubernetesExecutor与LocalKubernetesExecutor这类静态编码混合执行器并非核心 Airflow 的原生机制而是通过任务实例的queue字段来标记并持久化该任务应该跑在哪个子执行器上。这滥用了queue字段导致在使用这些混合执行器时无法再用队列做它本来的用途如 Celery worker 的队列分流、队列限流等。组合爆炸问题每增加一种执行器组合都需要手工编写一个新的具体类例如 CeleryKubernetes、LocalKubernetes……随着执行器种类增多这种手工组合不可持续维护成本越来越高。因此从 Airflow 3.0.0 开始CeleryKubernetesExecutor不再受支持官方推荐统一使用多执行器并发特性让用户以声明式配置自由组合任意执行器。迁移路径使用多执行器并发Using Multiple Executors Concurrently从Airflow 2.10.0开始Airflow 原生支持多执行器并发配置这也是官方指定的CeleryKubernetesExecutor替代方案。迁移后不再需要特化类只需在配置里列出一个或多个执行器即可。配置方式沿用单执行器的[core] executor配置项用逗号分隔指定多个执行器executor/index.rst[core] executor CeleryExecutor,KubernetesExecutor要点列表中第一个执行器是环境的默认执行器未显式指定执行器的任务 / DAG 都由它执行其行为与 2.10.0 之前完全一致列表中的其他执行器会被初始化并随时待命只有出现在列表中的执行器才可用于运行任务自定义执行器也可以用完整模块路径例如executor KubernetesExecutor,my.custom.module.ExecutorClass。别名Alias为便于在 DAG 中引用配置支持为执行器指定别名executor/index.rst[core] executor my_local_exec:LocalExecutor,my_celery_exec:CeleryExecutor也可以混用内置名与别名甚至结合多租户Multi-Team用法按团队隔离[core] executor global_celery_exec:CeleryExecutor;team1team_celery_exec:CeleryExecutor注意同一个团队内不能使用同一 Executor 类的两个实例例如一个团队不能用两个 CeleryExecutor 实例该限制仅在多租户模式下放宽。详见核心文档 multi-team.rst。在 DAG 与任务上指定执行器任务级指定executor/index.rstfrom airflow.operators.bash import BashOperator BashOperator( task_idhello_world, executorKubernetesExecutor, # 该任务强制走 Kubernetes bash_commandecho hello world!, )装饰器写法from airflow.decorators import task task(executorCeleryExecutor) def hello_world(): print(hello world!)DAG 级指定——通过default_args让整个 DAG 的任务默认使用某个执行器可被单个任务覆盖with DAG( dag_idhello_worlds, default_args{executor: CeleryExecutor}, # 应用到 DAG 内所有任务 ) as dag: hw hello_world() hw_again hello_world_again()任务在 DAG 解析时把所选执行器持久化到 Airflow 元数据库中后续重新解析 DAG 会同步更新。迁移对照与监控差异能力CeleryKubernetesExecutor2.x已废弃多执行器并发2.10推荐选择执行器的依据任务queue kubernetes_queue默认kubernetes任务 / DAG 上的executor参数组合方式需要预制的混合类配置中逗号分隔任意组合队列字段语义被劫持用于路由无法再用于队列分流队列字段恢复本意路由与队列解耦支持范围Airflow 3.0 起移除Airflow 2.10 与 3.x 持续支持监控方面单执行器时指标行为与 2.9 及之前一致配置多个执行器后executor.open_slots、executor.queued_slots、executor.running_tasks等指标会按执行器分别发布指标名追加执行器类名例如executor.open_slots.CeleryExecutor便于分别观测两条执行链路的负载executor/index.rst。日志行为与单执行器场景一致。总结CeleryKubernetesExecutor是 Airflow 演进史上的一个阶段性方案它通过按队列路由在同一调度器内同时驱动 Celery 与 Kubernetes 两套执行引擎让海量轻量任务享受 Celery 的吞吐、让少量重型任务享受 Kubernetes 的隔离。但其对queue字段的占用以及每组合一个类的实现方式注定只能作为过渡形态存在。如果你仍在使用 Airflow 2.x 并运行该执行器建议尽快规划迁移升级到 2.10 后改用[core] executor CeleryExecutor,KubernetesExecutor的多执行器并发配置通过任务 / DAG 级executor参数实现等价甚至更灵活的分流并利用按执行器拆分的监控指标持续观测两条执行链路。当前仓库中的核心实现celery_kubernetes_executor.py、配置默认值provider_config_fallback_defaults.cfg与测试用例test_celery_kubernetes_executor.py可以为你理解路由细节、核对行为提供最直接的参考。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询