ᕕ( ᐛ )ᕗ Jimyag's Blog

队列已经空了,正在处理的任务能安全结束吗?

Last modified:

4 个 worker 各自领取一个耗时 60 秒的任务后,队列里的待领取任务数变成了 0。HPA 看到零积压,立即把副本从 4 缩到 1。被删除的 3 个 Pod 手里的任务怎么办?

队列已经空了,是否等于所有任务都完成了?

本文用 224 个真实任务和一次执行中缩容实验回答这个问题。它是 HPA 实验系列的第三篇;前两篇分别介绍了 CPU 控制循环应用吞吐指标

为什么生产系统会按队列积压扩缩容

1
2
3
用户提交任务 → broker pending → worker claim → inflight → ACK → completed
                       └→ task_queue_depth → HPA → worker 副本数

图片处理、转码、报表生成等异步任务常被外部 I/O 限制,CPU 可能不高,但等待任务仍不断增加。此时积压比 CPU 更直接地表达「还有多少工作没有开始」。

本例的 broker 为教学用内存实现,支持任务 ID、claim、ACK、租约和重试;consumer 每次只处理一个任务。生产环境还需要根据实际消息系统确认持久化、确认和重投语义。

Service 为什么会出现在指标配置里

基础实验把全局队列长度关联到普通 Kubernetes Service queue

1
2
Prometheus → queue Service:8080 → queue Pod 的 /metrics
HPA → custom.metrics.k8s.io → Adapter → Prometheus 中属于 Service/queue 的指标

Service 同时承担网络入口和指标对象标识,但两件事分别配置。HPA 不会直接访问 Service 的 HTTP 端口。

Adapter 用 Prometheus 标签建立映射:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
# 将 Prometheus 队列序列映射为 Service 指标。
- seriesQuery: 'demo_queue_depth{namespace!="",service!=""}'
  resources:
    overrides:
      namespace: {resource: namespace}
      service: {resource: service}
  name:
    matches: '^demo_queue_depth$'
    as: queue_depth
  # 本实验只有一个全局 exporter,取最大值避免重复计数。
  metricsQuery: 'max(<<.Series>>{<<.LabelMatchers>>}) by (<<.GroupBy>>)'

HPA 中则有两个不同对象:

 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
15
16
17
# 读取 queue Service 的总量,扩缩 consumer Deployment。
scaleTargetRef:
  apiVersion: apps/v1
  kind: Deployment
  name: consumer # 真正被扩缩的对象
metrics:
  - type: Object
    object:
      describedObject:
        apiVersion: v1
        kind: Service
        name: queue # 指标关联的对象
      metric:
        name: task_queue_depth
      target:
        type: AverageValue
        averageValue: "5" # 每个 consumer 分摊 5 个待领取任务

可以把它读成一句话:

查询 queue Service 关联的待领取任务数,按每个 consumer 5 个任务计算,调整 consumer Deployment。

Object + AverageValue 的原始建议是:

1
desiredReplicas = ceil(pendingTasks / 5)

describedObject 不是扩缩目标,也不会让 HPA 调用队列 API。直接查询指标 API,可以先确认映射是否正确:

1
2
kubectl --context kind-hpa-lab get --raw \
  '/apis/custom.metrics.k8s.io/v1beta1/namespaces/hpa-demo/services/queue/task_queue_depth'

基础链路中队列总量为 35 时,API 返回 value: "35"。4 个 worker 稳定后,HPA 将同一个总量显示成每副本平均值:

1
2
3
4
custom metrics API: value=35
HPA current.averageValue: 8750m
HPA desiredReplicas: 4
Deployment Ready: 4

这里的 8750m8.75 个任务,不是毫秒。原始 API 保留总量,AverageValue 的分摊发生在 HPA 的计算语义中。

先验证:更多 worker 能否加快处理

实验文件固定在 k8sdev c48f140。按第一篇创建集群后运行:

1
2
3
4
cd ../advanced
./setup.sh
HPA_OUTPUT=$(mktemp -d "${TMPDIR:-/tmp}/hpa-queue.XXXXXX")
python3 run.py real --output "$HPA_OUTPUT"

两组实验各提交 80 个任务,每个任务处理 1 秒:

方式 首次运行 9 月 8 日重跑 Ready 副本范围
固定单 consumer 80.36 秒 85.30 秒 1
HPA 驱动 32.10 秒 36.20 秒 1~4

固定单副本与 HPA 处理相同任务批次的完成进度和 Ready 副本数

图表来自首次运行;重跑从 enqueue 到最后一个 ACK 重新计时,并再次 PASS。完成数来自逐任务 ACK;副本数约每 5 秒采样一次。这个结果只适用于内存队列和 sleep 任务,不能推导真实业务会线性提速。

持续入队实验每秒提交一个耗时 2 秒的任务。单 consumer 理想处理能力约为 0.5 个/秒,先出现积压,随后 HPA 增加 consumer;停止生产后,pending 和 inflight 都归零,最终缩回 1 个副本。

队列目标如何选择

设入队速率为 λ,每个 worker 平均处理速率为 μ,有效副本数为 N

1
积压变化率 ≈ λ - N × μ

只有 N × μ > λ,历史积压才会持续减少。若希望在 D 秒内清掉现有 Q 个任务,可用下面的简化式估算起点:

1
N ≥ (λ + Q/D) / μ

本例的 AverageValue: 5 只表达每副本分摊 5 个 pending 任务,没有表达完成期限。任务从 1 秒变成 60 秒时,队列数量相同,所需处理能力却不同。生产阈值应使用实际耗时分布、可接受等待时间和下游并发上限校准。

关键实验:pending=0 时触发缩容

实验先维持 4 个 consumer,再提交 4 个各耗时 60 秒的任务。确认 inflight=4 后,将最小副本改回 1、缩容稳定窗口设为 0。

1
2
3
pending=0, inflight=4
HPA: desiredReplicas 4 → 1
3 Pods: Running → Terminating

一个被缩容的 consumer 留下了完整时间线:

1
2
3
4
14:32:13.599 claim id=221 duration_ms=60000
14:32:25.112 SIGTERM: stop claiming; wait for current task
14:33:13.603 ack id=221 draining=true
14:33:13.603 drained; exit

另外两个被缩容的 Pod 对任务 222、223 也保持相同顺序;剩余 Pod 完成任务 224。最终审计结果:

字段 结果
accepted / completed 224 / 224
pending / inflight 0 / 0
retries 0
每个任务 attempts 1

程序在退出时停止 claim,等待当前任务结束,发送 ACK 后退出;Pod 配置 terminationGracePeriodSeconds: 100,给最长 60 秒的任务留下收尾时间。完整依据在 queue-audit.json 和同目录的 consumer 日志。

为什么宽限期仍不等于 exactly-once

任务守恒关系是:

1
accepted = pending + inflight + completed

因此 pending=0 只能说明没有待领取任务。readiness 失败也只能影响 Service 流量,不会自动停止主动轮询 broker 的 consumer。

正常 SIGTERM 路径通过,不代表进程崩溃、节点断电或 ACK 丢失也安全。如果业务操作成功但 ACK 丢失,租约到期后的重投可能再次执行同一个操作。生产 consumer 通常仍需按任务 ID 做幂等,或把去重记录与业务更新放进合适的事务边界。

回答开头的问题

队列为空不等于任务完成:pending=0 时,任务可能都在 inflight。本次 3 个被缩容的 Pod 能安全结束,是因为应用在 SIGTERM 后停止领取、等待任务、发送 ACK,并且 Kubernetes 的终止宽限期足够长。

所以,队列积压适合决定需要多少 worker,但不能负责任务可靠性。缩容安全最终取决于 claim、ACK、租约、重试和幂等语义。 下一篇会沿着指标到 Ready Pod 的完整链路,定位 HPA 看起来「没有按预期工作」的原因。

#Kubernetes #HPA #消息队列 #Kind