ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Airflow CeleryExecutor 深度指南:集群扩展、任务分发与运维实践

Apache Airflow CeleryExecutor 深度指南:集群扩展、任务分发与运维实践 Apache Airflow CeleryExecutor 深度指南集群扩展、任务分发与运维实践【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow导读CeleryExecutor是 Apache Airflow 中将任务执行水平扩展到多台工作机Worker的核心执行器方案。本文以 Airflow Celery Provider 的官方文档为骨架结合仓库内 Celery Executor 的真实源码celery_executor.py、default_celery.py与单元测试系统讲解 CeleryExecutor 的部署前提、airflow celery全套 CLI 命令、Broker/Result Backend 选型、Redis 维护规范、架构与任务执行时序、结构化日志以及队列Queue路由机制。读完本文你将能够独立规划并运维一套基于 Celery 的 Airflow 分布式任务集群。一、CeleryExecutor 是什么为什么需要它CeleryExecutor是 Airflow 在生产环境中将任务执行规模横向扩展的途径之一调度器Scheduler只负责解析 DAG、决定哪些任务实例需要运行而真正的任务执行被派发到由 Celery 管理的分布式 Worker 集群中。要让CeleryExecutor工作起来需要完成三件事准备一个 Celery BackendBroker常见选择包括RabbitMQ、Redis、Redis Sentinel等安装依赖包例如librabbitmq、redis等与所选 Broker 对应的客户端库修改airflow.cfg将[core] executor参数指向CeleryExecutor并配置相关 Celery 参数。Broker 的完整安装与搭建方法可参考 Celery 官方入门文档Airflow 侧的全部配置项集中在 Celery Provider 的 configurations-ref 中本文后续章节会结合源码逐一展开。Worker 节点的硬性要求文档明确列出了部署 Worker 时四条必须满足的前提条件airflow必须已安装且 CLI 位于 PATH 中Worker 进程依赖airflowCLI 来启动本地任务作业进程集群内 Airflow 配置必须保持一致homogeneous所有 Worker、Scheduler 应使用相同的airflow.cfg配置任务所用 Operator 的依赖必须在对应 Worker 上满足例如使用HiveOperator的机器必须装有 Hive CLI使用MySqlOperator则要求相应的 Python 库可通过PYTHONPATH被导入Worker 必须能访问DAGS_FOLDER文件系统需要自行同步常见做法是把DAGS_FOLDER放进 Git 仓库通过 Chef、Puppet、Ansible 等配置管理工具同步到各机器若所有机器共享同一挂载点将流水线文件放在共享目录同样可行。依赖不一致是 Celery 集群中最常见的“任务在某台机器失败、换一台就成功”问题的根源务必在部署阶段用统一的镜像或配置管理工具固化 Worker 环境。二、启动与停止 Workerairflow celery命令族启动 Worker配置就绪后在目标机器上启动 Workerairflow celery workerWorker 启动后会立即开始监听队列一旦有任务被投递到它订阅的队列就会拉取并执行。该命令的底层实现位于 celery_command.py 的worker()函数中其组装出的 Celery 实际启动参数为worker -O fair --queues 默认队列 --concurrency worker_concurrency --loglevel 级别-O fair启用 Celery 的 fair 调度策略避免长任务阻塞后续任务--queues默认取[operators] default_queue配置--concurrency默认取[celery] worker_concurrency默认值 16见 get_provider_info.py。worker子命令还支持以下常用参数定义见 definition.py参数说明-q, --queues逗号分隔的队列名列表worker 只处理这些队列的任务-c, --concurrencyworker 进程数任务实例并发数-H, --celery-hostname一台机器上运行多个 worker 时指定主机名-a, --autoscale形如max,min的动态扩缩容启用后worker_concurrency被忽略-t, --team多租户multi-team模式下的团队名需要 Airflow 3.2 且开启[core] multi_team-u, --umask守护模式下 worker 的 umask--without-mingle/--without-gossip启动时不与其他 worker 同步 / 不订阅其他 worker 事件停止 Workerairflow celery stop该命令读取 worker 的 PID 文件向 Celery 主进程发送SIGTERM信号实现优雅停机graceful shutdown——正在执行的任务会跑完收尾流程随后移除 PID 文件。这与 Celery 官方推荐的停止方式一致。更多运维命令除worker、stop外CLI 还提供了一批实用的 worker 运维命令全部位于 celery_command.pyairflow celery list-workers列出所有活跃 worker 及其监听的队列支持-o table|json|yaml|plain输出格式airflow celery shutdown-worker -H celeryhostname请求指定 worker 优雅关闭airflow celery shutdown-all-workers -y向所有活跃 worker 广播关闭带二次确认airflow celery add-queue -H celeryhostname -q queue1,queue2运行时为 worker 动态订阅队列airflow celery remove-queue -H celeryhostname -q queue1运行时取消订阅队列airflow celery remove-all-queues -H celeryhostname清空某 worker 的所有队列订阅。注意shutdown-worker、add-queue、remove-queue等远程控制类命令依赖 Celery 的远程控制remote control能力若在配置中将[celery] worker_enable_remote_control False这些命令将不可用且 Flower 也无法工作。三、监控Celery Flower 与安装方式你可以运行Celery Flower——一个构建在 Celery 之上的 Web UI用于监控 worker 状态、队列长度与任务执行情况airflow celery flower使用前必须确保flowerPython 库已安装。官方推荐直接安装 Airflow 的 Celery 组合包bundlepip install apache-airflow[celery]从仓库 README.rst 可以看到Celery Provider 的依赖为celery[redis] 5.5.0,6与flower 1.0.0。flower子命令支持以下选项-H/--hostname默认取[celery] flower_host默认0.0.0.0、-p/--port默认取flower_port默认5555、-u/--url-prefixURL 前缀、-A/--basic-auth形如user1:password1,user2:password2的基础认证、-a/--broker-api与-c/--flower-conf。其实现会组装出flower --addresshost --portport [--broker-api...]等参数交给 Celery app 启动。四、部署注意事项Caveats官方文档给出了一份浓缩的“避坑清单”逐条说明如下务必使用数据库支撑database-backed的 Result Backend不要用 Redis/AMQP 等键值型后端保存任务结果原因见“Redis 维护”章节在[celery_broker_transport_options]中设置visibility_timeout且必须超过最长任务的 ETA否则长任务会被 Broker 判定超时并重新投递造成重复执行使用 Redis Sentinel 作为 Broker、且 Redis 服务端有密码保护时必须在[celery_broker_transport_options]中通过sentinel_kwargs指定密码需符合 JSON 字典格式如{password: password_for_redis_server}设置[celery] worker_umask控制 Worker 新建文件的权限位为 Worker 配置充足资源worker_concurrency决定了单机同时执行的任务数资源不足会导致任务互相拖垮队列名称长度限制为 256 字符且不同 Broker 后端可能还有各自的额外限制。从源码看visibility_timeout的默认行为在 default_celery.py 的_broker_transport_options()中可以看到当配置中未显式设置visibility_timeout时若 Broker URL 以redis://、rediss://、sqs://、sentinel://开头Airflow 会**自动填入 86400 秒24 小时**并打印警告日志单元测试 test_celery_executor.py 验证了这一行为未配置时默认 86400 且产生警告显式配置后无警告对不支持该参数的 Broker如 RabbitMQ则不设置。此外[celery] task_acks_late默认True会在任务完成后才向 Broker 确认但对 Redis/SQS Brokertask_acks_late无法覆盖visibility_timeout——任务执行超过visibility_timeout仍会被重投递。因此长任务场景下两者必须配合调整。五、Redis 维护规范重要当 Redis 作为 Celery Broker 时需要明确 Redis 与 Airflow 元数据库的职责边界Redis 只承载瞬态 Celery 消息排队中的任务命令、任务确认acknowledgement数据以及其他 Broker 侧状态Airflow 任务历史、Dag Run、连接Connections、变量Variables与 XCom全部存储在 Airflow 元数据数据库中应使用数据库工具如airflow db clean维护。因此清理元数据库不会清 Redis清理 Redis 也不会清元数据库两者必须分别维护。维护窗口的正确姿势禁止在 Scheduler 或 Celery Worker 运行期间FLUSHDB/FLUSHALL或删除 Airflow 所使用的 Redis 键——排队中、执行中或正在确认的任务可能因此丢失或使任务状态短暂不一致。若停机后确需丢弃过期 Broker 数据请按以下顺序操作先停止 Airflow Schedulers 与 Celery Workers确认没有需要保留的queued/running状态任务若 Redis 开启持久化先备份 Redis 数据库仅删除[celery] broker_url配置所指向的 Redis 数据库或 keyspace。若你同时把 Redis 配置为 Celery 的result_backend维护窗口还需要一并处理 result backend 的键——删除这些键会移除 Scheduler 仍可能查询的 Celery 任务结果信息。生产环境官方强烈建议使用数据库支撑的 Result Backend使 Broker 清理与任务结果存储彻底分离。六、架构全景组件与通信CeleryExecutor下的 Airflow 集群由以下组件构成对应架构图中的编号通信关系组件职责Workers执行被分配的任务Scheduler把需要执行的任务放入队列Web serverHTTP 服务提供 DAG/任务状态查询Database存放任务状态、DAG、Variables、Connections 等元数据Celery队列机制由Broker与Result backend两部分组成Celery 队列的两个组成部分分工明确Broker存储待执行的命令任务消息Result backend存储已完成命令的状态。组件间 11 条通信路径编号通信说明1Web server → Workers拉取任务执行日志2Web server → Dag files展示 DAG 结构3Web server → Database拉取任务状态4Workers → Dag files读取 DAG 结构并执行任务5Workers → Database读写连接配置、变量、XCom6Workers → Celery result backend保存任务状态7Workers → Celery broker存取执行命令8Scheduler → Dag files读取 DAG 结构9Scheduler → Database存储 Dag Run 与相关任务10Scheduler → Celery result backend获取已完成任务的状态11Scheduler → Celery broker投放待执行命令七、任务执行时序从调度到结果回写图CeleryExecutor 下从任务入队到状态回写的完整时序源自 Celery Provider 官方文档。参与进程与存储任务执行前已存在的常驻进程与存储SchedulerProcess处理任务运行于 CeleryExecutor 之上WorkerProcess监听队列等待新任务出现WorkerChildProcess等待被分配新任务QueueBroker任务消息队列ResultBackend任务结果存储。执行过程中会临时创建两个进程LocalTaskJobProcess其逻辑由LocalTaskJob描述负责监控 RawTaskProcess通过TaskRunner启动新进程RawTaskProcess承载用户代码的进程例如BaseOperator.execute方法。13 步执行流程[1]SchedulerProcess 处理任务发现需要执行的任务后将其发送到QueueBroker[2]SchedulerProcess 开始周期性地向ResultBackend查询任务状态[3]QueueBroker 感知到任务后将任务信息发送给某个WorkerProcess[4]WorkerProcess 将单个任务分配给一个WorkerChildProcess[5]WorkerChildProcess 执行 Celery 任务处理函数对应 celery_executor.py 中的execute_command并创建新进程LocalTaskJobProcess[6]LocalTaskJobProcess 的逻辑由LocalTaskJob类描述使用TaskRunner启动新进程[7][8]任务完成后RawTaskProcess与LocalTaskJobProcess被停止[10][12]WorkerChildProcess 向主进程WorkerProcess通知任务结束以及可接收后续任务[11]WorkerProcess 将状态信息保存到ResultBackend[13]SchedulerProcess 再次向 ResultBackend 查询时获得任务最终状态。源码中的同步机制在调度器侧CeleryExecutorcelery_executor.py负责以下关键工作批量投递_send_workloads_to_celery()使用ProcessPoolExecutor多进程并行调用send_workload_to_executor投递任务单条或单进程时退化为主线程直接发送并发度由[celery] sync_parallelism控制默认0表示使用max(1, CPU核数-1)个进程状态回写update_all_workload_states()通过BulkStateFetcher批量查询AsyncResult把 Celery 的SUCCESS/FAILURE/REVOKED/STARTED/PENDING/RETRY状态映射为 Airflow 任务状态success/fail任务接管try_adopt_task_instances()支持 Scheduler 重启或 HA 切换后依据预分配的外部执行器 IDexternal_executor_id即 Celerytask_id接管此前提交的任务避免重复执行任务撤销revoke_task()通过celery_app.control.revoke(task_id)撤销排队中的任务发布重试任务发布遇AirflowTaskTimeout时会按[celery] task_publish_max_retries默认 3 次重试期间会递增celery.task_timeout_error统计指标。八、Worker 日志从纯文本到结构化 JSON默认情况下Celery Worker 将纯文本日志写入 stdout。若要输出结构化 JSON可通过[logging]段全局开启作用于所有组件或仅对 Worker 用[celery]段单独覆盖# 全局生效——影响 API server、scheduler 和所有 worker [logging] json_logs True # 仅覆盖 Celery worker其他组件保持不变 [celery] json_logs True配置查找顺序[celery] json_logs—— 若设置则优先采用[logging] json_logs——[celery]段缺失时回退使用False—— 两处均未配置时的默认值。这套回退逻辑与 Worker 启动代码中已有的[logging] CELERY_LOGGING_LEVEL→[logging] LOGGING_LEVEL回退机制保持一致见 celery_command.py 中worker()对日志级别的读取。在 celery_command.py 的实现中json_logs在 Airflow 3.1 会传入configure_logging(json_output...)以切换 structlog 输出Airflow 3.0.x 则因 SDK 版本限制而忽略该参数。版本注意[logging] json_logs是 Airflow 3.2.0 新增的全局键。在更早的 3.x 版本上该全局键会被静默忽略fallbackFalse。因此只设置[celery] json_logs True是在所有核心版本上都安全的开启方式。另外若设置[logging] celery_stdout_stderr_separation TrueWorker 日志会按级别分流ERROR 及以上写入 stderr其余写入 stdoutlogger_setup_handler回调实现。九、队列Queues按需路由任务CeleryExecutor允许为任务指定 Celery 队列。queue是BaseOperator的一个属性因此任意任务都可以分配到任意队列。环境默认队列由airflow.cfg的[operators] default_queue定义它既是任务未显式指定队列时的归属队列也是 Worker 启动时监听的默认队列Worker 可监听一个或多个队列启动时通过逗号分隔无空格传入队列名集合例如airflow celery worker -q spark,quark该 Worker 此后只处理被路由到spark或quark队列的任务。在CeleryExecutor的配置构建中task_default_queue与task_default_exchange均取自[operators] default_queue见 default_celery.py。队列的应用场景队列机制特别适合以下两类需求资源维度例如非常轻量的任务单个 Worker 就能承载数千个可为其单独开一条队列环境维度例如希望 Worker 运行在 Spark 集群内部因为它需要非常特定的运行环境与安全权限。任务侧通过queuespark之类的算子参数指定队列配合airflow celery worker -q与运行时动态add-queue/remove-queue命令即可实现灵活的任务路由。十、附录Celery 核心配置速查以下参数来自 get_provider_info.py 中的 Provider 配置声明即 configurations-ref 的数据来源按配置段归纳[celery]段配置项默认值说明celery_app_nameairflow.providers.celery.executors.celery_executorCelery app 名称worker_concurrency16airflow celery worker启动的并发进程数worker_autoscale未设置形如16,12的最大/最小动态扩缩容worker_prefetch_multiplier1Worker 预取倍数1 可能造成任务被长任务阻塞worker_enable_remote_controlTrue是否启用 Worker 远程控制关闭后 Flower 不可用worker_umask未设置守护模式下 Worker 的文件创建掩码八进制mp_start_method未设置stdlibmultiprocessing启动方式fork/forkserver/spawnbroker_urlredis://redis:6379/0Celery Broker URL敏感项result_backend未设置未指定时自动使用dbsql_alchemy_conn敏感项result_backend_sqlalchemy_engine_optionsResult backend SQLAlchemy 引擎选项如{pool_recycle: 1800}flower_host/flower_port/flower_url_prefix/flower_basic_auth0.0.0.0/5555//Flower 监听地址、端口、URL 前缀、基础认证sync_parallelism0状态同步进程数0表示max(1, 核数-1)ssl_active/ssl_mutual_tls/ssl_key/ssl_cert/ssl_cacertFalse/True///Broker SSL/TLS 配置mTLS 时ssl_key、ssl_cert必填poolpreforkCelery 池实现prefork、eventlet、gevent、solooperation_timeout1.0send_workload_to_executor/fetch_celery_task_state操作超时秒task_acks_lateTrue任务完成后才确认Redis/SQS 不覆盖 visibility_timeouttask_track_startedTrue任务开始执行时上报started状态支持 HA 下接管task_publish_max_retries3任务发布超时最大重试次数extra_celery_config{}额外 Celery 配置字典如{worker_max_tasks_per_child: 10}json_logs未设置Worker 单独启用 JSON 日志优先于[logging] json_logs[celery_broker_transport_options]段配置项说明visibility_timeout消息可见性超时秒未设置时 Redis/SQS 默认 86400必须大于最长任务执行时长sentinel_kwargsRedis Sentinel 客户端参数JSON 字典格式如{password: ...}敏感项[celery_result_backend_transport_options]段配置项说明master_nameRedis Sentinel 作为 result backend 时必填的主节点名sentinel_kwargsResult backend 的 Sentinel 客户端参数JSON 字典格式敏感项client-config/fetch_message_attributes/predefined_exchanges/predefined_queues/queue_tags/sqs-creation-attributesSQS 相关选项JSON 字典格式kafka_common_config/kafka_admin_config/kafka_consumer_config/kafka_producer_configKafka 相关选项JSON 字典格式以上所有“JSON 字典格式”选项在 default_celery.py 中会被json.loads解析为真实字典后传给 Celery若格式非法将直接抛出ValueError提示“should be written in the correct JSON format”。Result Backend 的自动推导在 default_celery.py 的get_default_celery_config()中若未显式配置result_backendAirflow 会自动拼接元数据库连接将[database] sql_alchemy_conn加上db前缀如dbpostgresql://...并依据环境中的 SQLAlchemy/psycopg 版本自动选择postgresqlpsycopg://或postgresqlpsycopg2://方言。同时若检测到 result backend 使用redis://、amqp://、rpc://等非数据库协议会打印“强烈建议改用数据库作为 result backend”的告警——这与前文“务必使用数据库支撑的 result backend”的运维建议完全对应。关于 Airflow 与 Python 的模块管理细节Worker 上依赖如何被发现可进一步参考 modules_management。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表