Apache Airflow Celery Provider 配置参考:CeleryExecutor 全部配置项、默认值与源码级解析

发布时间:2026/9/14 11:05:27
Apache Airflow Celery Provider 配置参考:CeleryExecutor 全部配置项、默认值与源码级解析 Apache Airflow Celery Provider 配置参考CeleryExecutor 全部配置项、默认值与源码级解析【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowApache Airflow 的 Celery providerapache-airflow-providers-celery通过CeleryExecutor将任务分发到独立的 Celery worker 集群是大规模水平扩展 Airflow 的核心方案之一。本篇基于仓库中的providers/celery/docs/configurations-ref.rst配置参考页及其数据来源 provider.yaml 完整展开逐节讲解[celery]、[celery_broker_transport_options]、[celery_result_backend_transport_options]、[celery_kubernetes_executor]四个配置段的全部参数、类型、默认值与敏感项规则并结合 default_celery.py 的源码说明这些配置最终如何被解析、组装进 Celery 应用读完即可在生产环境中正确配置 broker、结果后端、TLS 与 worker 行为。配置参考页的生成机制providers/celery/docs/configurations-ref.rst 本身并不直接罗列参数而是通过 Sphinx include 复用通用模板其实际参数内容来自两处数据源通用模板 devel-common/src/sphinx_exts/includes/providers-configurations-ref.rst 提供页面标题与说明本页包含该 provider 所有可配置的 Airflow 配置项可通过airflow.cfg文件或环境变量设置参数清单模板 devel-common/src/sphinx_exts/includes/sections-and-options.rst 遍历configs字典为每个 section 和 option 渲染描述、类型、默认值、示例与环境变量名。而configs数据的权威定义在 provider.yaml 的config:键中最新 provider 版本为 3.24.0。构建文档时会将其同步到自动生成的 get_provider_info.py——该文件头部注释明确写着THIS FILE IS AUTOMATICALLY GENERATED。因此下文所有参数的描述、类型、默认值均可在这两个文件中交叉验证。渲染规则来自 sections-and-options.rst 模板还包括每个 section 渲染为[section_name]标题选项按字典序列出带version_added的选项标注引入版本如task_acks_late标注 3.6.0一批 SQS/Kafka 传输选项标注 3.24.0标记sensitive: true的选项额外提供_CMD与_SECRET后缀的环境变量见下文。配置段总览该 provider 共提供四个配置段provider.yaml 中的config键配置段作用前提条件[celery]CeleryExecutor的全部行为参数airflow.cfg的[core] executor CeleryExecutor[celery_broker_transport_options]传给底层 Celery broker transport 的选项同上[celery_result_backend_transport_options]传给结果后端 transport 的选项Redis Sentinel、SQS、Kafka 等同上[celery_kubernetes_executor]CeleryKubernetesExecutor的队列路由参数[core] executor CeleryKubernetesExecutor[celery]与[celery_broker_transport_options]两个段在 provider.yaml 中的描述均注明only applies if you are using the CeleryExecutor in[core]section即这些配置只在对应 executor 启用时才生效。[celery] 配置段应用与 worker 基础参数参数类型默认值说明celery_app_namestringairflow.providers.celery.executors.celery_executorCelery 使用的 app 名称worker_concurrencystring16airflow celery worker启动 worker 时的并发度即单个 worker 同时处理的 task instance 数。应依据 worker 机器资源与任务性质调整worker_autoscalestring未设置形如max_concurrency,min_concurrency的自动扩缩容区间示例16,12。一旦设置 autoscaleworker_concurrency将被忽略worker_prefetch_multiplierinteger1worker 预取任务数 进程数 × 该值。大于 1 可提升吞吐但多 worker 场景下可能让已被其他 worker 认领的长任务阻塞后续任务worker_enable_remote_controlbooleantrue是否启用 worker 远程控制。broker 不支持远程控制时 Celery 会创建大量.*reply-celery-pidbox队列可设为false避免但关闭后 Flower 无法工作poolstringpreforkCelery Pool 实现可选prefork默认、eventlet、gevent、soloworker_umaskstring未设置daemon 模式下 worker 的 umask八进制整数决定 worker 新建文件的初始权限位mp_start_methodstring未设置airflow celery worker进程为日志服务器、stale-bundle 清理进程及可选的[secrets] use_cache管理器等标准库multiprocessing辅助进程使用的启动方法取值须为当前平台multiprocessing.get_all_start_methods()返回之一通常fork、forkserver、spawn。未设置时回退到[core] mp_start_method再回退到平台默认mp_forkserver_preloadstring未设置逗号分隔的模块列表forkserver进程启动时预先导入使 worker 的multiprocessing辅助进程通过写时复制继承而非重新导入仅在有效mp_start_method为forkserver时使用未设置时回退到[core] mp_forkserver_preload示例airflow关于mp_start_methodprovider.yaml 中给出了关键背景Python 3.14 将 Unix 默认启动方法从fork改为forkserver后者以及spawn会在每个辅助进程中重新导入 Airflow 并额外拉起 forkserver/resource-tracker 进程从而增加 worker 常驻内存设置mp_start_method fork可恢复 3.14 之前的行为。同时该设置只约束标准库multiprocessing辅助进程Celery 的prefork池由billiardmultiprocessing的独立分支驱动、始终使用fork不受此配置影响。broker 与结果后端参数类型敏感默认值说明broker_urlstring是redis://redis:6379/0Celery broker 地址支持多种 broker 类型RabbitMQ、Redis、Redis Sentinel、SQS 等result_backendstring是未设置Celery 结果后端。任务完成后需更新元数据状态供 scheduler 查询强烈建议使用数据库后端未设置时自动使用[database] sql_alchemy_conn加db前缀。示例dbpostgresqlpsycopg2://postgres:airflowpostgres/airflowresult_backend_sqlalchemy_engine_optionsstring否传给 Celery 结果后端 SQLAlchemy 引擎的可选配置字典JSON 字符串示例{pool_recycle: 1800}result_backend的描述中还包含一个版本约束示例中的psycopg2驱动在所有受支持的 Airflow 版本上均可用只有 Airflow 3.2.0 及以后才应使用dbpostgresqlpsycopg://因为更早版本不保证 SQLAlchemy 2.0而 psycopg (v3) 方言存在于该版本线中。这段逻辑在 default_celery.py 中得到印证get_default_celery_config()先检查team_conf.has_option(celery, RESULT_BACKEND)未配置时读取SQL_ALCHEMY_CONN根据运行时是否同时具备 SQLAlchemy 2.0 与psycopg包模块顶部的_USE_PSYCOPG3探测选择postgresqlpsycopg://或postgresqlpsycopg2://前缀替换postgresql://最后拼装为db{sql_alchemy_conn}。此外源码中还有两处值得注意的行为若result_backend以rediss?://、amqp://或rpc://开头仅发出日志告警提示建议改用数据库后端match_not_recommended_backend逻辑并不阻止启动除显式配置外源码还会向 Celery 配置中写入若干与[operators]段联动的固定项task_default_queue与task_default_exchange取自[operators] DEFAULT_QUEUEaccept_content/event_serializer固定为jsonAirflow 3 上还会默认设置worker_redirect_stdoutsFalse与worker_hijack_root_loggerFalse。调度器侧同步参数参数类型默认值说明sync_parallelismstring0CeleryExecutor同步任务状态所用进程数0表示使用max(1, 核心数 - 1)operation_timeoutfloat1.0send_workload_to_executor或fetch_celery_task_state操作的超时秒数task_publish_max_retriesinteger3发布任务消息到 broker 因AirflowTaskTimeout失败时的最大重试次数超出后放弃并标记任务失败celery_config_optionsstringairflow.providers.celery.executors.default_celery.DEFAULT_CELERY_CONFIGCelery 配置选项的导入路径即上文源码实现的入口extra_celery_configstring{}额外 Celery 配置JSON 字典字符串worker 启动时生效可覆盖任意 Celery 配置项例如{worker_max_tasks_per_child: 10}extra_celery_config在源码中的生效位置见 default_celery.py通过team_conf.getjson(celery, extra_celery_config, fallback{})解析后以字典展开**合并进最终配置且排在基础键之后因此可以覆盖前面由配置推导出的同名字段——这是任何 celery config 都能加进来的实现依据。任务语义参数参数类型默认值说明task_acks_lateboolean3.6.0 引入True任务执行时间超过visibility_timeout时broker 会把任务重新分配给其他 worker即使原任务仍在正常跑造成同一任务并发执行Airflow UI/日志只能看到Task Instance Not Running错误。设为True让任务完成后再 ack。注意对 Redis 与 SQS brokertask_acks_late并不能覆盖visibility_timeoutbroker 仍会按时重投长任务必须同时调大[celery_broker_transport_options] visibility_timeout默认 86400 秒task_track_startedbooleanTrueworker 执行任务时上报started状态供 Airflow 跟踪运行中任务scheduler 重启或 HA 模式下可借此收养adopt前任 SchedulerJob 遗留的孤儿任务json_logsboolean未设置以 JSON 格式输出 Celery worker 的 stdout。设置后优先于全局[logging] json_logs让 worker 与部署中其他组件使用不同日志格式未设置时回退到[logging] json_logs其自身默认Falsetask_track_started与task_acks_late的取值在 default_celery.py 中以fallbackTrue形式写入 Celery 配置即 provider 层面的兜底默认值。TLS 参数组参数类型默认值说明ssl_activestringFalse是否为 broker 连接启用 SSLssl_mutual_tlsbooleanTrue是否要求双向 TLS客户端证书认证True时ssl_key与ssl_cert必须同时设置False则为单向 TLS仅验证服务端ssl_keystring客户端私钥路径ssl_mutual_tls True时必需ssl_certstring客户端证书路径ssl_mutual_tls True时必需ssl_cacertstringCA 证书路径未设置时使用系统 CA 信任库验证服务端default_celery.py 的 SSL 处理逻辑可以读出几条硬性规则ssl_active解析失败如配置了非法布尔值时不会崩溃而是打警告日志will run without SSL并回退为关闭ssl_mutual_tls True默认但ssl_key/ssl_cert缺失时抛出ValueError明确提示要么两者都配要么把 SSL_MUTUAL_TLS 设为 False仅支持amqps?://RabbitMQ与rediss?:///sentinel://Redis/Sentinel两类 broker URLRabbitMQ 走cert_reqs/ca_certs/keyfile/certfileRedis 走ssl_cert_reqs/ssl_ca_certs/ssl_keyfile/ssl_certfile字段其他 URL 直接报错请改用 RabbitMQ 或 Redisssl_mutual_tls False却配置了 key/cert 时只告警并忽略客户端证书。Flower 监控参数组参数类型敏感默认值说明flower_hoststring否0.0.0.0FlowerCelery 的 Web 监控 UI绑定 IP由airflow celery flower命令读取flower_portstring否5555Flower 监听端口flower_url_prefixstring否Flower 的根 URL 前缀示例/flowerflower_basic_authstring是为 Flower 启用 Basic 认证接受逗号分隔的user:password对示例user1:password1,user2:password2这组参数直接被 CLI 定义消费cli/definition.py 中ARG_FLOWER_HOSTNAME、ARG_FLOWER_PORT、ARG_FLOWER_URL_PREFIX、ARG_FLOWER_BASIC_AUTH的默认值分别来自conf.get(celery, FLOWER_HOST)、conf.getint(celery, FLOWER_PORT)等即配置文件中的这四个 key 就是airflow celery flower -H/-p/-u/-A的默认来源。[celery_broker_transport_options] 配置段该段用于向底层 Celery broker transport 透传选项Celery 的broker_transport_options。当前定义了 2 个参数参数类型敏感默认值示例说明visibility_timeoutstring否未设置21600等待 worker ack 任务的秒数超时后消息重投给其他 worker。未设置时Airflow 对 Redis 与 SQS broker 默认取 86400 秒24 小时超时未 ack 的任务会被终止并重投。务必将其设为大于最长任务的运行时间。仅支持 Redis 与 SQS brokersentinel_kwargsstring是未设置{password: password_for_redis_server}传给 Redis Sentinel 客户端的附加选项Redis Sentinel 作 broker 且服务器设了密码时必须经此传入。类型虽是 string但必须是符合字典格式的 JSON 字符串源码印证见 default_celery.py_broker_supports_visibility_timeout()判定redis://、rediss://、sqs://、sentinel://四类 URL 支持可见性超时仅当用户未显式设置时才对上述 broker 注入默认86400并打出 warning 提示调大该值sentinel_kwargs与另外 10 个选项client-config、kafka_*_config、predefined_queues等被归入_BROKER_TRANSPORT_DICT_OPTIONS列表统一做 JSON 解析解析失败抛出should be written in the correct JSON format的ValueError。[celery_result_backend_transport_options] 配置段该段用于向结果后端的 transport 透传选项在 Redis Sentinel 作结果后端、SQS/Kafka 等场景下尤为重要。当前定义了 12 个参数其中 3.24.0 新增一批 SQS 与 Kafka 选项参数版本敏感示例说明master_name—否mymaster使用 Redis Sentinel 作结果后端时需连接的 Sentinel 主节点名sentinel_kwargs—是{password: password_for_redis_server}传给结果后端 Sentinel 客户端的附加选项要求字典格式 JSON 字符串client-config3.24.0否{connect_timeout: 5}SQS 的 botocore 客户端配置fetch_message_attributes3.24.0否{MessageSystemAttributeNames: [SenderId, SentTimestamp]}SQS 拉取消息时携带的属性predefined_exchanges3.24.0否{exchange-1: {arn: arn:aws:sns:us-east-1:xxx:exchange-1}}SQS 预定义 SNS 主题predefined_queues3.24.0否{queue-1: {url: https://sqs.us-east-1.amazonaws.com/xxx/aaa}}SQS 预定义队列queue_tags3.24.0否{Environment: production, Team: backend}建队时应用的 SQS 队列标签sqs-creation-attributes3.24.0否{KmsMasterKeyId: alias/aws/sqs}建队时附加的 SQS 队列属性键值对kafka_common_config3.24.0否{bootstrap.servers: broker:9094}Kafka producer/consumer 公共配置kafka_admin_config3.24.0否—Kafka 管理端专用配置kafka_consumer_config3.24.0否—Kafka 消费者专用配置kafka_producer_config3.24.0否—Kafka 生产者专用配置所有参数类型均为 string字典项要求 JSON 字符串。sentinel_kwargs的 JSON 解析逻辑在 default_celery.py 中实现解析失败会抛出AirflowException提示应使用正确的字典格式。[celery_kubernetes_executor] 配置段参数类型默认值说明kubernetes_queuestringkubernetes使用CeleryKubernetesExecutor时任务所在队列等于该值则由KubernetesExecutor执行否则走CeleryExecutorCeleryKubernetesExecutor是一个包装器celery_kubernetes_executor.py 中kubernetes_queue是cached_property通过conf.get(celery_kubernetes_executor, kubernetes_queue)读取本参数_router()方法则据此在queue_command/queue_task_instance等入口把任务分发给内部两个 executor。需要注意源码中的兼容性限制在 Airflow 3.0 上直接实例化CeleryKubernetesExecutor会抛出RuntimeError提示改用并发使用多个 executor的新机制见其__init__中的错误信息。环境变量的设置规则配置参考页为每个选项渲染了对应的环境变量命名规则为AIRFLOW__{SECTION}__{OPTION}section 与 option 均转大写、点号换下划线普通选项AIRFLOW__CELERY__BROKER_URL、AIRFLOW__CELERY__WORKER_CONCURRENCY、AIRFLOW__CELERY_BROKER_TRANSPORT_OPTIONS__VISIBILITY_TIMEOUT、AIRFLOW__CELERY_RESULT_BACKEND_TRANSPORT_OPTIONS__SENTINEL_KWARGS等标记sensitive的选项broker_url、result_backend、flower_basic_auth、两处sentinel_kwargs额外提供三个变量便于从命令输出或密钥服务注入AIRFLOW__CELERY__BROKER_URLAIRFLOW__CELERY__BROKER_URL_CMD执行该命令的标准输出作为取值AIRFLOW__CELERY__BROKER_URL_SECRET从指定文件中读取取值一个可复制的最小配置示例综合以上各段一份典型的airflow.cfg配置如下仅展示本 provider 相关部分[core] executor CeleryExecutor [celery] # broker支持 RabbitMQ / Redis / Redis Sentinel / SQS 等 broker_url redis://redis:6379/0 # 结果后端强烈建议用数据库不配置则自动取 sql_alchemy_conn 加 db 前缀 result_backend dbpostgresqlpsycopg2://postgres:airflowpostgres/airflow worker_concurrency 8 # 长任务必须保证可见性超时大于最长任务运行时间 task_acks_late True # 额外覆盖任意 Celery 原生配置 extra_celery_config {worker_max_tasks_per_child: 10} [celery_broker_transport_options] visibility_timeout 86400 # Redis Sentinel 且服务端有密码时必须设置 # sentinel_kwargs {password: password_for_redis_server} [celery_result_backend_transport_options] # 仅在使用 Redis Sentinel 作结果后端时设置 # master_name mymaster # sentinel_kwargs {password: password_for_redis_server}对应的 worker 启动/停止命令由 cli/definition.py 注册的 CLI 提供# 启动 worker-q 指定监听的队列默认取 [operators] default_queue # -c 指定并发默认取 worker_concurrency-a 指定 max,min 自动扩缩 airflow celery worker -q spark,quark -c 8 # 优雅停止向主 Celery 进程发送 SIGTERM airflow celery stop # 启动 Flower 监控参数默认值来自上文 flower_* 配置项 airflow celery flower -p 5555 -u /flower该 CLI 组实际包含 9 个子命令worker、flower、stop、list-workers、shutdown-worker、shutdown-all-workers、add-queue、remove-queue、remove-all-queues其中队列订阅类命令add/remove-queue、shutdown-worker要求-H传入完整 celery 主机名如celeryhostname。相关文档与源码入口配置参考页本体providers/celery/docs/configurations-ref.rst参数权威数据源providers/celery/provider.yamlconfig:键provider 3.24.0生成后的 Python 表示get_provider_info.pyCelery 配置组装实现default_celery.pyget_default_celery_config即celery_config_options默认导入路径指向的对象CeleryExecutor 使用指南worker 要求、架构、队列、日志celery_executor.rstCLI 命令定义cli/definition.py混合 executor 实现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个关键决策

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

获取专属建站方案

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

立即免费咨询