Flower 功能全景指南:Celery 集群的实时监控、远程控制与可观测性集成 可观测性运维后端【免费下载链接】flowerReal-time monitor and web admin for Celery distributed task queue项目地址https://gitcode.com/gh_mirrors/fl/flower点击查看免费下载Flower 是 Celery 分布式任务队列的实时监控与 Web 管理工具。本文围绕其核心功能清单展开以 Celery Events 为基础的实时监控、基于 Celery Control 的远程管理worker 关停/扩容/限速/撤销任务、Broker 队列监控以及 HTTP Basic Auth 与 Google/GitHub/GitLab/Okta OAuth 认证、Prometheus 指标导出和 HTTP API 这五大能力域逐一拆解其操作方式、底层实现与配置要点。读完本文你将掌握 Flower 的核心功能边界、各功能的调用链与配置方法能据此搭建一套完整的 Celery 可视化监控与治理体系。功能总览Flower 的五大能力域Flower 面向 Celery 集群提供两条互补的管线被动观测通过订阅 Celery Events 事件流实时掌握 worker 在线状态、任务全生命周期状态与历史主动控制借助 Celery 的 remote control远程控制命令通道对 worker 执行关停、扩容、改队列、设限、撤销任务等管理操作。这两条管线加上Broker 监控、认证体系HTTP Basic Auth 与多种 OAuth、Prometheus 集成和REST API构成了 Flower 的完整功能地图。仓库中的功能清单文档 docs/features.rst 将其概括如下基于 Celery Events 的实时监控任务进度与历史、任务参数/启动时间/运行时长等详情远程控制worker 状态与统计、关停与重启、池大小与自动扩缩容、队列增删、运行中/定时/预留/撤销任务查看、时间与速率限制、任务撤销与终止Broker 监控所有 Celery 队列的统计信息HTTP Basic Auth 与 Google/GitHub/GitLab/Okta OAuthPrometheus 集成HTTP API下文逐项展开。实时监控基于 Celery Events 的任务与 Worker 观测事件流的采集机制Flower 的监控数据源是 Celery 的 Events 机制。worker 在运行任务过程中会向 broker 发送task-sent、task-received、task-started、task-succeeded、task-failed、task-retried、worker-online、worker-heartbeat、worker-offline等事件Flower 通过 broker 连接订阅这些事件并聚合到内存状态中。核心实现在 flower/events.pyEvents是一个后台守护线程threading.Thread通过EventReceiver持续capture事件flower/events.py每收到一个事件都会通过io_loop.add_callback回到 Tornado 的 ioloop 线程调用EventsState.event避免跨线程同步问题flower/events.pyEventsState继承自 Celery 的celery.events.state.State并在其上维护每个 worker 的事件计数器self.counterflower/events.py。Flower 还内置了周期性发送enable_events命令的机制on_enable_events每 5 秒执行一次capp.control.enable_events确保在 Flower 启动之后才上线的 worker 也能开启事件上报flower/events.py 与 flower/events.py。因此即使你启动 worker 时没有加-E参数Flower 也能补开事件。提示--enable_events选项默认值为True见 flower/options.py。若 worker 端想自行开启可在启动 worker 时加-E标志。任务监控视图进度、历史与详情在 Flower 的 Web UI 中/tasks页面对应 flower/views/tasks.py 的TasksView以 DataTable 形式列出任务每个任务行展示任务名、UUID、状态、参数、关键字参数、结果、接收/启动时间、运行时长、worker 等字段可通过--tasks-columns自定义列。数据由TasksDataTable通过iter_tasks从事件状态中分页、排序、搜索得到flower/views/tasks.py。点击单个任务进入/task/task_id详情页TaskView可查看任务完整信息包括args/kwargs任务入参received/started/succeeded接收、启动、完成时间戳runtime实际运行时长result任务结果state当前状态PENDING/RECEIVED/STARTED/SUCCESS/FAILURE/RETRY 等worker执行该任务的 workereta/expires/retries/exception/traceback等扩展信息。iter_tasks的实现支持按type、worker、state、received_start/received_end、started_start/started_end过滤并支持sort_by含-前缀的降序与分页flower/utils/tasks.py。任务搜索语法Flower 支持 GitHub 风格的搜索语法详见 docs/tasks_filter.rst在任务列表的搜索框中即可使用搜索词含义foo查找 args、kwargs 或 result 中包含foo的所有任务args:foo查找参数中包含foo的任务kwargs:foobar查找关键字参数中foobar的任务result:foo查找结果中包含foo的任务state:FAILURE查找所有失败任务如果搜索词含空格需要用双引号包裹例如args:hello world。底层解析与匹配逻辑位于 flower/utils/search.pyparse_search_terms负责拆分查询词支持引号内的空格satisfies_search_terms负责逐项匹配——kwargs匹配通过对字符串化的字典做快速索引查找实现避免为每个任务反序列化字典flower/utils/search.py。状态持久化可选Flower 默认把全部 worker 与任务状态保存在内存中。若希望重启 Flower 后状态不丢失可以开启持久化模式--persistentTrue启用持久化--dbfile指定状态数据库文件默认flower--state-save-intervalms周期性保存间隔毫秒默认0表示不自动保存。对应实现中Events在持久化模式下用shelve保存/恢复整个状态对象并通过PeriodicCallback按间隔触发save_stateflower/events.py 与 flower/events.py。内存上限由--max-workers默认 5000与--max-tasks默认 100000控制见 flower/options.py。远程控制通过 Celery Control 管理 Worker 与任务Flower 的远程控制能力全部基于 Celery 的controlremote control commands通道实现并包装为 Web 页面按钮与 REST API。控制器实现集中在 flower/api/control.py路由注册见 flower/urls.py。Worker 生命周期管理关停 workerPOST /api/worker/shutdown/worker底层调用capp.control.broadcast(shutdown, destination[workername])flower/api/control.py。注意这会让 worker 进程优雅退出需由 supervisor/systemd 等守护进程拉起。重启 worker 池POST /api/worker/pool/restart/worker发送pool_restart命令reloadFalse。需要 worker 以--pool-restarts即CELERYD_POOL_RESTARTS启动才可用否则返回 403flower/api/control.py。池大小与自动扩缩容扩容POST /api/worker/pool/grow/worker?n3调用capp.control.pool_grow(nn, replyTrue, destination[...])n默认 1flower/api/control.py缩容POST /api/worker/pool/shrink/worker?n2调用capp.control.pool_shrinkflower/api/control.py自动扩缩容POST /api/worker/pool/autoscale/worker?min3max10发送autoscale广播命令。需要 worker 使用--autoscale启动对应CELERYD_AUTOSCALER否则返回 403flower/api/control.py。这些操作对应 UI 上 worker 详情页flower/templates/worker.html的Pool区域按钮。单元测试对每类操作均有覆盖例如test_pool_grow断言pool_grow(n3, replyTrue, destination[test])被正确调用而test_shutdown_read_only验证只读模式下所有控制操作返回 403 且不触发底层命令tests/unit/api/test_control.py。队列消费管理增加队列消费POST /api/worker/queue/add-consumer/worker?queuesample-queue发送add_consumer广播取消队列消费POST /api/worker/queue/cancel-consumer/worker?queuesample-queue发送cancel_consumer广播。两者实现见 flower/api/control.py。这使你可以动态调整某 worker 实例监听的队列而无需重启 worker。任务状态查看当前运行任务通过 inspect 命令active安全模式获取定时任务ETA/countdown通过 inspect 命令scheduled获取预留与撤销任务分别通过 inspect 命令reserved与revoked获取。Inspector将stats、active_queues、registered、scheduled、active、reserved、revoked、conf八个 inspect 方法并行提交到线程池执行并将结果按 worker 聚合flower/inspector.py。其中active使用safeTrue以携带任务参数信息。worker 列表页/workersWorkersView正是基于这些结果渲染每个 worker 的状态、统计与任务摘要flower/views/workers.py。时间限制与速率限制设置时间限制POST /api/task/timeout/taskname表单参数soft、hard、workername调用capp.control.time_limit(taskname, hard..., soft..., destination...)。若不指定workername则作用于所有 workerflower/api/control.py设置速率限制POST /api/task/rate-limit/taskname表单参数ratelimit如200/m、workername调用capp.control.rate_limit(...)flower/api/control.py。撤销与终止任务POST /api/task/revoke/task_id?terminatetruesignalSIGTERM未加terminate时只撤销尚未开始的任务加terminatetrue会向正在运行该任务的进程发送signal默认SIGTERM强制终止。底层调用capp.control.revoke(taskid, terminate..., signal...)flower/api/control.py。只读模式以上所有控制操作都受read_only选项约束开启--read_only后UI 隐藏操作按钮API 层面每个 handler 都会先检查self.application.options.read_only并返回 403例如 flower/api/control.py。该选项定义见 flower/options.py非常适合只做展示、不让误操作的运维场景。Broker 监控Celery 队列统计Flower 的/broker页面BrokerView展示所有 Celery 队列的统计信息如队列消息数。其实现根据 broker 类型选择不同的客户端RabbitMQ通过 RabbitMQ Management HTTP API 查询队列信息需要设置broker_api选项如http://guest:guestlocalhost:15672/api/。Flower 依据 broker URL 自动构造默认 API 地址也可显式覆盖flower/utils/broker.pyRedis / Redis Sentinel / Redis 单机套接字 / Redis SSL通过 Redis 的LLEN命令按优先级步长汇总队列长度flower/utils/broker.py。Broker 客户端的选择由工厂类Broker根据 broker URL 的 scheme 决定amqp/amqps→ RabbitMQredis→ Redisrediss→ RedisSslredissocket→ RedisSocketsentinel→ RedisSentinelflower/utils/broker.py。API 层面GET /api/queues/length会返回所有活跃队列的名称与消息数flower/api/tasks.py示例响应{ active_queues: [ {name: celery, messages: 0}, {name: video-queue, messages: 5} ] }配置broker_api的两种方式# 命令行 celery flower --broker-apihttp://username:passwordrabbitmq-server-name:15672/api/ # flowerconfig.py broker_api http://guest:guestlocalhost:15672/api/注意RabbitMQ Management Plugin 默认未启用需先执行rabbitmq-plugins enable rabbitmq_managementRabbitMQ 3.0 之前的管理端口是 55672。认证体系Basic Auth 与 OAuthFlower 支持多种认证方式详见 docs/auth.rst其中/healthcheck与/metrics两个端点始终豁免认证。HTTP Basic Auth最简单的方式适合内网快速加固celery flower --basic-authuser:pswd支持逗号分隔的多个username:password对例如--basic-authuser1:password1,user2:password2。选项定义见 flower/options.py。Google / GitHub / GitLab / Okta OAuthOAuth 认证由auth_provider选择具体 handler定义于 flower/views/auth.py配合auth允许访问的邮箱正则、oauth2_key、oauth2_secret、oauth2_redirect_uri四个选项使用。auth支持的基础正则语法包括单邮箱userexample.com通配符.*example.com管道分隔的邮箱列表oneexample.com|twoexample.com出于安全考虑auth仅支持上述基础语法其合法性校验在validate_auth_option中完成flower/views/auth.py。Google OAuth示例配置写入flowerconfig.pyauth_provider flower.views.auth.GoogleAuth2LoginHandler auth allowed-emails.*gmail.com oauth2_key your_client_id oauth2_secret your_client_secret oauth2_redirect_uri http://localhost:5555/login需要在 Google Developer Console 创建 OAuth 2.0 Client IDApplication type 选择 Web application并把http://localhost:5555/login加入 Authorized redirect URIs。GitHub OAuth将auth_provider改为flower.views.auth.GithubLoginHandler其余选项相同需先在 GitHub Settings 注册 OAuth App。源码中 GitHub 登录域名可用环境变量FLOWER_GITHUB_OAUTH_DOMAIN覆盖默认github.comflower/views/auth.py。GitLab OAuth将auth_provider设为flower.views.auth.GitLabLoginHandleroauth2_key/oauth2_secret取 GitLab 应用注册返回的 Application ID 与 Secret。可选用环境变量进一步限定FLOWER_GITLAB_AUTH_ALLOWED_GROUPS逗号分隔的允许组列表子组用/表示如group1,group2/subgroupFLOWER_GITLAB_MIN_ACCESS_LEVEL最低组访问级别FLOWER_GITLAB_OAUTH_DOMAIN自定义 GitLab 域名。Okta OAuth将auth_provider设为flower.views.auth.OktaLoginHandler同时通过FLOWER_OAUTH2_OKTA_BASE_URL环境变量提供 Okta 基础 URL。认证成功的用户邮箱会被写入安全 cookie后续请求自动通过未授权邮箱会收到 403 拒绝flower/views/auth.py。Prometheus 集成导出与可视化指标导出端点Flower 安装后即自带/metrics端点默认localhost:5555/metrics以 Prometheus 文本格式输出指标flower/views/monitor.py。指标在 flower/events.py 中定义完整列表如下指标名描述Labels类型flower_events_totalFlower 注册的 Celery 任务事件次数task, type, workercounterflower_task_prefetch_time_seconds任务在 worker 中等待被执行的时间task, workergaugeflower_worker_prefetched_tasksworker 上预取的某类任务数量task, workergaugeflower_task_runtime_seconds任务运行耗时task, workerhistogramflower_worker_onlineworker 在线状态workergaugeflower_worker_number_of_currently_executing_tasksworker 当前执行中的任务数workergauge其中flower_task_runtime_seconds的直方图桶可用--task-runtime-metric-buckets1,5,10,inf自定义默认 PrometheusHistogram.DEFAULT_BUCKETS选项定义见 flower/options.py。worker-online/worker-heartbeat/worker-offline事件分别驱动flower_worker_online的 1/1/0 切换task-received事件递增预取任务计数task-started事件计算预取耗时并扣减计数flower/events.py 与 flower/events.py。配置 Prometheus 抓取在prometheus.yml的scrape_configs中把 Flower 加入抓取目标scrape_configs: - job_name: prometheus static_configs: - targets: [localhost:9090] - job_name: flower static_configs: - targets: [localhost:5555]仓库根目录自带一份可直接使用的示例 prometheus.yml其 job 名为flowertargets 为flower:5555。若要抓取完整的 Prometheus 指标必须保证 worker 开启了事件上报启动参数-E或依赖 Flower 的enable_events补开事件数据是这些指标的唯一来源。指标 Labels 的 PromQL 用法可借助 PromQL 按标签过滤task任务名如tasks.add、tasks.multiplytype任务事件类型如task-started、task-succeeded注意worker 相关事件不计入该指标workerworker 名如celeryhostname。例如统计失败任务总数sum(flower_events_total{typetask-failed}) by (task)告警规则示例仓库提供了一份现成的 Prometheus 告警规则 examples/prometheus-alerts.yaml包含三条典型规则- alert: CeleryWorkerOffline expr: flower_worker_online 0 for: 2m labels: severity: critical context: celery-worker annotations: summary: Celery worker offline description: Celery worker {{ $labels.worker }} has been offline for more than 2 minutes. - alert: TaskFailureRatioTooHigh expr: (sum(avg_over_time(flower_events_total{typetask-failed}[15m])) by (task) / sum(avg_over_time(flower_events_total{type~task-failed|task-succeeded}[15m])) by (task)) * 100 1 for: 5m - alert: TaskPrefetchTimeTooHigh expr: sum(avg_over_time(flower_task_prefetch_time_seconds[15m])) by (task, worker) 1 for: 5m将规则内容加入你的 Alertmanager 规则文件中即可生效。Grafana 监控看板仓库提供了开箱即用的 Grafana 看板 examples/celery-monitoring-grafana-dashboard.json。导入步骤在 Grafana 侧边栏点击→Import→Upload JSON file选择该文件 → 为看板指定一个 Prometheus 数据源 → 点击Import完成。导入后即可看到包含 worker 在线状态、任务执行量与失败率、任务运行耗时分布等面板的 Celery 监控看板作为自建看板的起点。完整的本地化演练路径Redis broker Celery app Flower Prometheus Grafana见 docs/prometheus-integration.rst核心启动命令如下# 1. 启动 Redis broker docker run --name redis -d -p 6379:6379 redis # 2. 配置 Celeryexamples/celeryconfig.py # broker_url redis://localhost:6379/0 # celery_result_backend redis://localhost:6379/0 # 3. 启动 worker-E 开启事件 celery -A tasks worker -l INFO -E # 4. 启动 Flower celery -A tasks --brokerredis://localhost:6379/0 flower # 5. 启动 Prometheus挂载包含 flower job 的 prometheus.yml docker run --name Prometheus -v ABSOLUTE PATH TO prometheus.yml:/etc/prometheus/prometheus.yml \ -p 9090:9090 --network host prom/prometheus # 6. 启动 Grafana docker run --name Grafana -d -v grafana-storage:/var/lib/grafana \ -p 3000:3000 --network host grafana/grafana提示仓库根目录的 docker-compose.yml 与 prometheus.yml 可帮助你以容器方式快速拉起整套环境examples/tasks.pyexamples/tasks.py提供了add、sleep、echo、error等示例任务其中error任务专门用于制造失败场景便于验证告警。HTTP API程序化访问 Flower 能力Flower 将监控与控制能力全部暴露为 REST API路由清单见 flower/urls.py方法路径功能GET/api/workers?refreshstatusworkername列出 worker 及其状态/统计/注册任务/配置POST/api/worker/shutdown/worker关停 workerPOST/api/worker/pool/restart/worker重启 worker 池POST/api/worker/pool/grow/worker?n扩容池POST/api/worker/pool/shrink/worker?n缩容池POST/api/worker/pool/autoscale/worker?minmax自动扩缩容POST/api/worker/queue/add-consumer/worker?queue增加队列消费POST/api/worker/queue/cancel-consumer/worker?queue取消队列消费GET/api/tasks?limitoffsetsort_byworkernametasknamestatereceived_startreceived_endsearch列表任务可搜索过滤GET/api/task/info/task_id获取任务详情POST/api/task/apply/taskname执行任务并等待结果POST/api/task/async-apply/taskname异步执行任务POST/api/task/send-task/taskname按名称发送任务无需任务源码GET/api/task/result/task_id?timeout查询任务结果POST/api/task/abort/task_id中止运行中的任务POST/api/task/timeout/taskname设置 soft/hard 时间限制POST/api/task/rate-limit/taskname设置速率限制POST/api/task/revoke/task_id?terminatesignal撤销/终止任务GET/api/queues/length所有活跃队列长度GET/api/task/types已观测到的任务类型GET/metricsPrometheus 指标GET/healthcheck健康检查返回 OK以任务结果查询为例curl http://localhost:5555/api/task/result/c60be250-fe52-48df-befb-ac66174076e6 # {result: 3, state: SUCCESS, task-id: c60be250-fe52-48df-befb-ac66174076e6}认证与只读约束下的 API 行为出于安全考虑API 默认在未启用认证时不可用BaseApiHandler.prepare会检查是否配置了basic_auth或auth两者都未配置且未显式开启FLOWER_UNAUTHENTICATED_API时返回 401flower/api/init.py。如需在无认证环境开放 API设置export FLOWER_UNAUTHENTICATED_APItrue注意这只应在可信内网使用。此外所有写操作关停、扩容、撤销、限速等在read_only模式下统一返回 403。任务执行类接口/api/task/apply通过IOLoop.run_in_executor将阻塞等待结果的操作放到线程池避免阻塞事件循环flower/api/tasks.py/api/task/abort依赖 Celery 的 Abortable 任务机制且需要配置结果后端否则 503flower/api/tasks.py。常用启动方式与配置入口Flower 作为 Celery 的子命令运行命令行模板为celery [celery options] flower [flower options]直接指定 brokercelery --brokerredis:// flower --unix-socket/tmp/flower.sock使用既有 Celery 应用celery -A tasks.app flower修改端口celery -A tasks.app flower --port5001默认 5555通过配置文件默认加载flowerconfig.py可用--conf覆盖文件为普通 Python 键值对# 设置 RabbitMQ management api broker_api http://guest:guestlocalhost:15672/api/ # 启用 debug 日志 logging DEBUG通过环境变量所有选项加FLOWER_前缀后均可通过环境变量传入例如export FLOWER_BASIC_AUTHfoo:bar celery flower命令行参数经过apply_options的解析流程先解析命令行获得--conf再加载配置文件最后以命令行覆盖环境变量则在apply_env_options中统一注入flower/command.py。常见选项及默认值汇总选项默认值说明port5555HTTP 服务端口address空串监听地址空串监听所有接口unix_socket改用 UNIX socket 监听basic_authNone逗号分隔的user:pass对auth允许访问的邮箱正则auth_providerNoneOAuth handler 类路径oauth2_key/oauth2_secret/oauth2_redirect_uriNoneOAuth 2.0 凭据与回调broker_apiNoneRabbitMQ Management API 地址inspect_timeout1000inspect 命令超时毫秒max_workers/max_tasks5000 / 100000内存中保留的 worker/任务上限db/persistent/state_save_intervalflower / False / 0状态持久化相关enable_eventsTrue周期性下发 enable_eventsnatural_timeFalse时间显示为相对时间tasks_columnsname,uuid,state,args,kwargs,result,received,started,runtime,worker/tasks页显示的列url_prefix部署在非根路径时的前缀cookie_secret随机 64 字符cookie 签名密钥read_onlyFalse禁用所有控制操作task_runtime_metric_bucketsPrometheus 默认桶任务运行时长直方图桶purge_offline_workersNone离线 worker 清理时间秒完整选项定义见 flower/options.py。若在反向代理后面部署需配置Host头且url_prefix必须与代理路径一致如FLOWER_URL_PREFIXflower配合location /flower/与proxy_pass http://localhost:5555/flower/否则前端静态资源链接会解析失败同时建议用 nginx 的auth_basic做额外访问控制docs/reverse-proxy.rst。小结Flower 通过事件订阅 远程控制 Broker 客户端 认证/指标/API的组合把 Celery 集群的观测与管理能力收敛到一个 Web 界面和一套 REST API 中实时监控依赖 Celery Events任务状态、历史与详情全部来自事件流并支持持久化与自定义搜索远程控制全部映射到 Celery control 通道从 worker 关停、池扩容到队列调整、限速与撤销任务均可编程触发Broker 监控通过 RabbitMQ Management API 或 Redis 命令直接读取队列深度可观测性通过/metrics直接对接 Prometheus仓库附带了告警规则与 Grafana 看板模板安全与集成覆盖 Basic Auth 与四种 OAuth 方案并在只读模式下可一键冻结全部写操作。如需进一步深入建议直接阅读仓库中的 flower/events.py、flower/api/control.py、flower/utils/broker.py 以及对应单元测试 tests/unit/api/test_control.py并结合 docs/config.rst 的完整选项参考按需调整部署参数。赞分享可观测性运维后端【免费下载链接】flowerReal-time monitor and web admin for Celery distributed task queue项目地址https://gitcode.com/gh_mirrors/fl/flower点击查看免费下载相关推荐Flower项目Celery集群监控与管理工具详解Flower项目Celery集群监控与管理工具详解 概述 Flower是一个基于Web的工具专门用于监控和管理Celery分布式任务队列集群。它为开发者提供可观测性运维后端Elasticsearch-js 监控与可观测性终极指南如何实时监控客户端性能和集群状态Elasticsearch js 监控与可观测性终极指南如何实时监控客户端性能和集群状态 Elasticsearch js 客户端提供了强大的 监控与可观测性后端搜索引擎MMPose 中的 RTMPose 实时全身姿态估计COCO-WholeBody 配置、模型库与 RTMW 训练实战MMPose 中的 RTMPose 实时全身姿态估计COCO WholeBody 配置、模型库与 RTMW 训练实战 本文基于 MMPose 仓库中 conf可观测性运维后端上一篇Guard未发布的V2特性预览CallerArgumentExpression与静态导入如何改变验证体验下一篇update-browserslist-db高级技巧如何自定义浏览器数据库更新策略创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考