ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Apache Airflow Executors 完全指南:社区 Provider 提供的执行器全景与配置实战

Apache Airflow Executors 完全指南:社区 Provider 提供的执行器全景与配置实战 Apache Airflow Executors 完全指南社区 Provider 提供的执行器全景与配置实战【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowExecutors 是 Airflow 中真正跑任务的机制调度器Scheduler把 Task Instance 交给 Executor由它决定任务是在本地进程执行、投递到消息队列由远端 Worker 拉取还是拉起 Kubernetes Pod / 云上容器来运行。本文以社区维护的 Provider 提供的全部 Executor 实现为主线结合 Airflow 核心的 BaseExecutor 接口与 ExecutorLoader 加载机制完整讲解每个 Executor 的定位、适用场景、配置方式以及如何编写自己的自定义 Executor帮助你在单机、多机、云原生等不同部署形态下做出正确的执行器选型。一、Executor 是什么Airflow 任务执行的插拔式引擎Airflow 官方核心文档 airflow-core/docs/core-concepts/executor/index.rst 给出了精确定义Executors are the mechanism by which task instances get run. They have a common API and are pluggable, meaning you can swap executors based on your installation needs.Executor 具有统一的公共 API 且可插拔你可以根据部署环境随时更换。一个容易被误解的点是Executor 的逻辑运行在 Scheduler 进程内部它只是决定任务在本地跑还是远程跑并不需要为 Executor 单独启动一个进程但像 Celery Worker 这类执行任务的远端组件仍需要单独部署。在 Airflow 源码中所有 Executor 的公共基类是BaseExecutor位于 airflow-core/src/airflow/executors/base_executor.py内置 Executor 的名字常量定义在 airflow-core/src/airflow/executors/executor_constants.py其中CORE_EXECUTOR_NAMES包含四个核心名字LOCAL_EXECUTOR LocalExecutor CELERY_EXECUTOR CeleryExecutor KUBERNETES_EXECUTOR KubernetesExecutor MOCK_EXECUTOR MockExecutor而ExecutorLoaderairflow-core/src/airflow/executors/executor_loader.py维护了核心名字到模块路径的映射executors { LOCAL_EXECUTOR: airflow.executors.local_executor.LocalExecutor, CELERY_EXECUTOR: airflow.providers.celery.executors.celery_executor.CeleryExecutor, KUBERNETES_EXECUTOR: airflow.providers.cncf.kubernetes.executors.kubernetes_executor.KubernetesExecutor, }从这里可以清楚看到CeleryExecutor 与 KubernetesExecutor 虽然名字上是核心执行器但其实现代码分别托管在 celery 与 cncf.kubernetes 两个社区 Provider 中。这正是关联文档Airflow can be extended by providers with Executors. Each provider can define their own Executors的直接体现——Provider 是 Airflow 扩展 Executor 的第一等公民机制。1.1 查看当前生效的 Executor通过 CLI 即可查询当前配置的 Executor$ airflow config get-value core executor LocalExecutor1.2 基本配置Executor 通过配置文件[core]段的executor选项设置这是官方文档明确的唯一配置入口[core] executor KubernetesExecutor自定义或第三方 Executor 则直接填写 Python 类的完整模块路径[core] executor my.custom.executor.module.ExecutorClass二、社区 Provider 提供的 Executor 全景关联文档指出社区维护的 Provider 提供了一批开箱即用的 Executor 实现。通过检索整个仓库中继承BaseExecutor的类grep class .*Executor(BaseExecutor)可以得到当前仓库中社区 Provider 提供的完整清单Executor 类所属 Provider源码路径定位CeleryExecutorapache-airflow-providers-celeryproviders/celery/src/airflow/providers/celery/executors/celery_executor.py基于 Celery 消息队列的分布式执行器KubernetesExecutorapache-airflow-providers-cncf-kubernetesproviders/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/kubernetes_executor.py每个任务一个 Kubernetes Pod 的容器化执行器EdgeExecutorapache-airflow-providers-edge3providers/edge3/src/airflow/providers/edge3/executors/edge_executor.py通过 HTTP 将任务分发到 Edge Worker 的远程执行器AwsBatchExecutorapache-airflow-providers-amazonproviders/amazon/src/airflow/providers/amazon/aws/executors/batch/batch_executor.py提交到 AWS Batch 运行AwsEcsExecutorapache-airflow-providers-amazonproviders/amazon/src/airflow/providers/amazon/aws/executors/ecs/ecs_executor.py提交到 AWS ECS 容器运行AwsLambdaExecutorapache-airflow-providers-amazonproviders/amazon/src/airflow/providers/amazon/aws/executors/aws_lambda/lambda_executor.py提交到 AWS Lambda 函数运行此外仓库中还保留了两种静态编码的混合 ExecutorExecutor 类所属 Provider源码路径CeleryKubernetesExecutorapache-airflow-providers-celeryproviders/celery/src/airflow/providers/celery/executors/celery_kubernetes_executor.pyLocalKubernetesExecutorapache-airflow-providers-cncf-kubernetesproviders/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/local_kubernetes_executor.py需要特别说明这两种混合 Executor 自 Airflow 3.0.0 起不再支持。官方核心文档明确解释了原因——它们的实现不是 Airflow 核心的原生能力而是通过滥用 Task Instance 的queue字段来指示和持久化由哪个子执行器运行这使得queue字段无法用于其本意同时每增加一种执行器组合都需要手工编写新的具体类维护成本随执行器数量增长而不可持续。官方推荐使用下文介绍的多 Executor 并发特性来替代。2.1 CeleryExecutor生产环境推荐的分布式执行器CeleryExecutor是社区中最成熟、最常用于生产环境的分布式执行器。其类文档明确指出CeleryExecutor is recommended for production use of Airflow. It allows distributing the execution of workloads (task instances and callbacks) to multiple worker nodes.从源码celery_executor.py可以看到它声明的能力标志supports_ad_hoc_ti_run: bool True supports_callbacks: bool True sentry_integration: str sentry_sdk.integrations.celery.CeleryIntegration pre_assigns_external_executor_id: ClassVar[bool] True supports_multi_team: bool True其核心工作流程对应_process_workloads/_send_workloads方法是调度器心跳期间把ExecuteTask或ExecuteCallback工作负载批量投递到 Celery broker如 Redis / RabbitMQ远端airflow celery worker从队列中拉取并执行。由于 Celery 不支持批量投递CeleryExecutor内部使用多进程池加速发送与状态查询相关配置项为celery.SYNC_PARALLELISM用于并行同步任务状态的进程数默认0表示自动取max(1, cpu_count() - 1)celery.task_publish_max_retries任务投递失败的最大重试次数默认3投递超时会累计重试并上报celery.task_timeout_error指标。启动 Celery Worker 的命令为源码 docstring 中明确给出airflow celery workerCelery 的bulk_state_fetcherBulkStateFetcher用于批量查询任务执行状态避免逐个查询造成的性能瓶颈。2.2 KubernetesExecutor任务级容器隔离KubernetesExecutor为每个 Task Instance 动态创建一个 Kubernetes Pod 来执行任务之间天然隔离无吵闹邻居问题。源码kubernetes_executor.py中同样声明了supports_ad_hoc_ti_run: bool True与supports_multi_team: bool True。其实现要点包括在__init__中通过KubeConfig(executor_confself.conf)读取 Kubernetes 相关配置并用self.kube_config.parallelism覆盖BaseExecutor的并行度使用multiprocessing.Manager()创建任务队列与结果队列惰性创建于start()中避免 API Server 构造执行器时泄漏 Manager 子进程内部AirflowKubernetesScheduler负责轮询 Pod 状态sync方法、失败重排队task_publish_max_retries对应配置项kubernetes_executor.task_publish_max_retries默认0通过get_task_log接口把运行任务的 Pod 日志注入 Airflow 任务日志便于排查执行环境自身导致的失败通过get_cli_commands提供airflow kubernetes generate-dag-yaml、airflow kubernetes cleanup-pods等运维命令。2.3 EdgeExecutorHTTP 分发的轻量远程执行器EdgeExecutoredge_executor.py的类文档描述为Implementation of the EdgeExecutor to distribute work to Edge Workers via HTTP.——它通过 HTTP 协议把工作负载分发给 Edge Worker适合在边缘/异构节点上执行任务的场景同样声明supports_multi_team: bool True。它的实现直接基于数据库表与 Worker 心跳start()中通过EdgeDBManager确保表结构存在任务状态通过EdgeJobModel、EdgeLogsModel、EdgeWorkerModel含EdgeWorkerState等模型持久化使用 SQLAlchemy 会话操作数据库并通过全局数据库锁DBLocks/create_global_lock保证并发安全。2.4 Amazon 系列对接云原生计算服务amazon Provider 的executors目录providers/amazon/src/airflow/providers/amazon/aws/executors/下集中了三个基于 AWS 计算服务的 ExecutorAwsBatchExecutor把任务作为 Job 提交到 AWS Batch 执行适用于批处理工作负载AwsEcsExecutor把任务作为 Task 提交到 ECS 集群的容器中执行每个任务一个容器实例AwsLambdaExecutor把任务作为 Lambda 函数调用执行适合轻量、短时任务。它们分别配套了batch_executor_config.py/ecs_executor_config.py任务级执行配置、boto_schema.pyBoto3 请求/响应模式以及utils/目录下的exponential_backoff_retry.py指数退避重试表明这三个执行器在重试、配置校验等工程细节上做了完整落地。三、Executor 的类型学本地 vs 远程核心文档将 Executor 分为两大类理解这一分类有助于选型。3.1 本地 Executor任务在 Scheduler 进程内直接以子进程方式运行。仓库中唯一的本地实现是LocalExecutorairflow-core/src/airflow/executors/local_executor.py也是 Airflow 的默认执行器。优点极易使用、速度快、延迟极低、几乎没有额外部署要求缺点能力受限且与 Scheduler 共享进程资源会影响调度器性能适用小型单机生产环境。3.2 远程 Executor远程执行器又细分为两类队列式 / 批处理式Queued/Batch Executors——任务投递到中央队列由常驻的远端 Worker 拉取执行典型如CeleryExecutor、AwsBatchExecutor、EdgeExecutor。优点Worker 与 Scheduler 解耦、更健壮Worker 可以是大型主机并行吞吐高、成本效率好Worker 常驻可即时拉取任务延迟相对较低缺点共享 Worker 存在吵闹邻居问题任务间争抢资源/环境配置若负载不恒定闲置或过度扩容的 Worker 会造成成本浪费。容器化Containerized Executors——任务排队时按需部署一个独立的容器/Pod 来执行如KubernetesExecutor、AwsEcsExecutor。优点每个任务独占容器环境无吵闹邻居问题执行环境可针对特定任务定制系统库、二进制、依赖、资源配额等Worker 仅在任务期间存活成本可控缺点容器/Pod 启动有延迟大量短小任务时开销偏高需要自行管理 Kubernetes 集群等基础设施。注意不要误以为远程 Executor 需要单独运行一个executor 进程。Executor 逻辑始终运行在 Scheduler 进程内它决定的是任务在本地还是远程执行。四、多 Executor 并发配置与别名从 Airflow 2.10.0 起支持多 Executor 并发每个 Executor 各有取舍延迟、隔离、计算效率之间的权衡让不同特性的任务跑在最适合它的执行器上可以扬长避短。4.1 逗号分隔的多执行器配置仍然使用[core] executor这一个配置项以逗号分隔多个执行器[core] executor LocalExecutor[core] executor LocalExecutor,CeleryExecutor[core] executor KubernetesExecutor,my.custom.module.ExecutorClass关键规则列表中的第一个执行器是环境的默认执行器任何未显式指定执行器的 Task 或 DAG 都会使用它其余执行器会被初始化并随时可用但只有被 DAG/Task 显式指定时才会运行任务。未列入配置的执行器无法被使用。4.2 别名Alias为了让 DAG 中指定执行器更简洁配置支持别名。别名对内置核心执行器和自定义模块路径都适用[core] executor LocalExecutor,short_name:my.custom.module.ExecutorClass[core] executor my_local_exec:LocalExecutor,my_celery_exec:CeleryExecutor在多团队Multi-TeamAirflow场景下别名尤其有用——同一个执行器类可以同时出现在全局与团队级别通过别名让任务精确指向某个实例[core] executor global_celery_exec:CeleryExecutor;team1team_celery_exec:CeleryExecutor限制说明同一个执行器类的两个实例仅在多团队模式下被支持例如两个团队可以各自使用 CeleryExecutor但单个团队内不允许配置两个 CeleryExecutor 实例一个执行器可以同时用于全局与团队。4.3 ExecutorLoader 的配置解析与校验逻辑从源码 executor_loader.py 的_get_team_executor_configs与_get_executor_names方法可以精确还原配置解析规则配置串按;分隔成多段每段形如团队名执行器列表或直接执行器列表无团队前缀表示全局配置每段内部按,拆分成单个执行器项每项按:拆分为别名:模块路径/核心名通过conf.get_mandatory_value(core, executor)强制要求该配置非空否则抛出AirflowConfigException校验规则源码中逐一检查并抛出带明确提示的异常至少存在一个全局执行器且全局执行器必须排在团队执行器之前同一团队内不允许重复配置同一个执行器按模块路径去重每个团队名只能出现一次团队执行器要求core.multi_team配置开启且团队名必须已存在于数据库中_validate_teams_exist_in_database只有supports_multi_teamTrue的执行器才能被配置为团队执行器别名形式的第二段必须是核心执行器名或包含.的合法模块路径。ExecutorLoader.load_executor支持四种加载方式核心执行器名如LocalExecutor、模块路径、执行器类名、或内部ExecutorName对象加载失败ImportError时会抛出提示检查[core] executor配置的AirflowConfigException。这些解析结果会被缓存_executor_names保证配置只解析一次。4.4 在 DAG 和 Task 中指定执行器Task 级别——通过 Operator 的executor参数BashOperator( task_idhello_world, executorLocalExecutor, bash_commandecho hello world!, )装饰器写法task(executorLocalExecutor) def hello_world(): print(hello world!)DAG 级别——通过default_args统一指定DAG 内所有任务默认使用该执行器单个任务仍可显式覆盖def hello_world(): print(hello world!) def hello_world_again(): print(hello world again!) with DAG( dag_idhello_worlds, default_args{executor: LocalExecutor}, # 应用于 DAG 内所有任务 ) as dag: hw hello_world() hw_again hello_world_again()重要提示如果 DAG 指定了未配置的执行器DAG 将解析失败并在 Airflow UI 中显示警告对话框——必须确保任何运行 Airflow 组件scheduler、workers 等的主机上配置都包含你想使用的全部执行器任务配置的执行器会持久化在 Airflow 数据库中每次 DAG 解析后更新。4.5 多执行器下的监控使用单执行器时指标行为与 2.9 及更早版本一致配置多个执行器后executor 指标executor.open_slots、executor.queued_slots、executor.running_tasks会按执行器分别发布并在指标名后追加执行器类名例如executor.open_slots.executor class name。日志行为与单执行器一致。五、编写自定义 ExecutorBaseExecutor 接口全解所有 Airflow Executor 都实现同一个公共接口调度器、日志、CLI 等组件都通过该接口与 Executor 交互这就是 airflow-core/src/airflow/executors/base_executor.py 中的BaseExecutor。需要自定义 Executor 的常见原因现有执行器不适合你的特定计算服务、需要使用云厂商的计算服务、或公司有私有的任务执行工具。5.1 工作负载WorkloadsWorkload 是 Executor 语境下的基本执行单元——一个离散的操作或作业例如在 Worker 上执行封装了用户代码的 Airflow 任务。仓库中工作负载类型定义在 airflow-core/src/airflow/executors/workloads/包括task.pyExecuteTask、callback.pyExecuteCallback、connection_test.pyTestConnection、trigger.py触发器、base.py与types.py。ExecuteTask的典型结构如下ExecuteTask( tokenmock, tiTaskInstanceDTO( idUUID(4d828a62-a417-4936-a7a6-2b3fabacecab), task_idmock, dag_idmock, run_idmock, try_number1, dag_version_idUUID(4d828a62-a417-4936-a7a6-2b3fabacecab), map_index-1, pool_slots1, queuedefault, priority_weight1, executor_configNone, ), dag_rel_pathPurePosixPath(mock.py), bundle_infoBundleInfo(namen/a, versionno matter), log_pathmock.log, typeExecuteTask, )5.2 BaseExecutor 的常用方法可不重写heartbeatScheduler Job 循环周期性调用是调度器与执行器的主要交互点负责更新指标、触发新任务执行、更新运行/完成任务状态queue_workload调度器通过它把工作负载交给执行器基类实现只是把 workload 加入内部待运行列表仓库中所有执行器都使用该方法get_event_buffer调度器调用以获取执行器当前执行中 TaskInstance 的状态has_task调度器判断执行器是否已持有某个任务的队列或运行状态send_callback把回调发送到执行器配置的 sink。5.3 必须实现的方法sync在心跳期间周期性调用更新执行器所知任务的状态可选地尝试执行已从调度器收到的排队任务execute_async异步执行一个工作负载通常只是把任务入队到内部或外部队列如KubernetesExecutor也可以直接执行如LocalExecutor_process_workloads处理通过queue_workload排队的负载列表定义执行器如何执行负载投递到 Worker、提交到外部系统等。5.4 可选实现的方法增强能力与稳定性start调度器初始化执行器对象后调用完成额外初始化如KubernetesExecutor惰性创建 Manager 队列end调度器拆除时调用做同步清理terminate更强硬地停止执行器甚至终止进行中的任务try_adopt_task_instances接管被遗弃的任务例如调度器崩溃遗留的无法接管的返回由基类处理get_cli_commands向airflowCLI 注入命令如 CeleryExecutor 的 worker 命令、KubernetesExecutor 的运维命令get_task_log向 Airflow 任务日志注入执行环境日志如 Kubernetes Pod 日志。get_cli_commands的伪代码模板staticmethod def get_cli_commands() - list[GroupCommand]: sub_commands [ ActionCommand( namecommand_name, helpDescription of what this specific command does, funclazy_load_command(path.to.python.function.for.command), args(), ), ] return [ GroupCommand( namemy_cool_executor, helpDescription of what this group of commands do, subcommandssub_commands, ), ]get_task_log的伪代码模板def get_task_log(self, ti: TaskInstance, try_number: int) - tuple[list[str], list[str]]: messages [] log [] try: res helper_function_to_fetch_logs_from_execution_env(ti, try_number) for line in res: log.append(remove_escape_codes(line.decode())) if log: messages.append(Found logs from execution environment!) except Exception as e: # 任何异常都不应导致任务日志失败 messages.append(fFailed to find logs from execution environment: {e}) return messages, [\n.join(log)]5.5 兼容性属性Compatibility AttributesBaseExecutor上的一组布尔属性供 Airflow 核心检查执行器的能力编写自定义执行器时必须正确设置属性含义supports_pickling是否支持在执行前从数据库读取 pickle 序列化的 DAG而非从文件系统读取 DAG 定义sentry_integration若支持 Sentry填写创建集成的可调用对象导入路径如CeleryExecutor设为sentry_sdk.integrations.celery.CeleryIntegrationis_local执行器是远程还是本地对应第三节的类型划分is_single_threaded是否单线程直接影响支持的数据库后端单线程执行器可运行于包括 SQLite 在内的任何后端is_production是否可用于生产非生产执行器会在 UI 中对用户显示提示serve_logs是否支持服务日志supports_multi_team是否支持多团队模式团队执行器必须为 True见上文校验逻辑5.6 接入自定义执行器实现BaseExecutor接口后把core.executor配置为模块路径即可生效[core] executor my_company.executors.MyCustomExecutor两个工程实践提示官方文档明确要求避免在模块顶层执行昂贵操作——Executor 类会在多处被导入导入缓慢会拖慢整个 Airflow 环境尤其是 CLI 启动从 Airflow 3.2.0 起可通过实现 provider 级 CLI 命令来管理核心扩展如执行器减少不必要的重量级导入缩短 CLI 启动时间。六、从源码看 Executor 的执行链路综合 base_executor.py 可以还原 Executor 与调度器的协作链路调度器启动后通过ExecutorLoader.init_executors()/get_default_executor()实例化配置的全部执行器并调用start()调度器循环周期性调用heartbeat()内部触发_process_workloads()把queue_workload收集的负载投递出去如 Celery broker、Kube API、HTTP Edge Worker与sync()同步任务状态执行器通过get_event_buffer()向调度器回报任务状态变化调度器拆除时调用end()/terminate()清理。值得注意的细节是base_executor.py中的get_execution_api_server_url它从api.base_url默认/推导出默认的执行 API 服务器地址http://localhost:8080/execution/并允许通过core.execution_api_server_url覆盖——这是 Airflow 3 中任务通过 HTTP 执行 API 上报状态这一新架构的基础设施多团队执行器可传入各自的ExecutorConf解析团队专属 URL。七、选型建议小结结合官方核心文档的利弊分析与上文源码证据可以给出以下务实的选择框架单机小规模 / 开发调试LocalExecutor零部署成本、低延迟多机分布式生产CeleryExecutor成熟稳定、社区推荐配合airflow celery worker弹性扩容任务级容器隔离 / 已有 K8s 集群KubernetesExecutor每个任务独占 Pod无吵闹邻居环境可定制边缘节点 / 异构 WorkerEdgeExecutor通过 HTTP 分发、基于数据库的轻量状态管理AWS 生态批处理选AwsBatchExecutor容器化选AwsEcsExecutor轻量短任务选AwsLambdaExecutor混合负载使用多 Executor 并发配置executor LocalExecutor,CeleryExecutor按任务特性分流。进一步深入可阅读的核心源码与文档执行器概念与完整配置airflow-core/docs/core-concepts/executor/index.rst基类接口airflow-core/src/airflow/executors/base_executor.py加载与校验逻辑airflow-core/src/airflow/executors/executor_loader.py本地执行器实现airflow-core/src/airflow/executors/local_executor.py各社区执行器源码providers/celery/src/airflow/providers/celery/executors/、providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/executors/、providers/edge3/src/airflow/providers/edge3/executors/、providers/amazon/src/airflow/providers/amazon/aws/executors/【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表