队列工作进程可观测性
队列深度表示有多少工作正在等待。队列等待说明了每个作业的积压成本。运行持续时间表示工作人员执行的时间。最终结果和重试次数表明工作最终是否成功。
您需要所有四个信号来区分健康的流量突发和落后的队列、缓慢的依赖项、重试风暴或消失的工作线程。
尾部百分位数揭示了平均水平可以消除的缓慢工作。
先决条件
- Telemetry API 钥匙
- 公开队列、开始和最终结果边界的工作人员
- 稳定的逻辑作业标识符和低基数作业名称
- 每个关键作业类型的预期完成窗口
分别对快照和结果进行建模
使用两种事件形状,因为它们回答不同的问题:
| 活动 | 谷物 | 使用 |
|---|---|---|
queue_snapshot |
一个采样时间一个队列 | 当前深度、可用工人、最长等待年龄 |
background_job_completed |
每个逻辑作业有一个最终结果 | 等待、运行持续时间、重试、成功、永久失败 |
快照是采样状态,因此不要随时间求和其队列深度。最终结果是已完成的工作,因此它无法找到已开始并消失的工作。当需要检测停滞作业时,添加轻量级 job_started 和 job_finished 生命周期事件。
定义终端事件契约
使用 后台作业完成架构 作为起点:
| 领域 | 含义 |
|---|---|
event_id |
用于重复数据删除的唯一终端事件 |
job_id |
跨尝试共享稳定的逻辑作业标识符 |
job_name |
有界作业类型,绝不是动态负载或 ID |
queue_name |
执行作业的队列 |
account_id |
受结果影响的假名账户 |
attempt_count |
包括终端尝试在内的总尝试次数 |
queue_wait_ms |
排队到第一次执行 |
duration_ms |
终端尝试执行持续时间 |
status |
success 或永久 error |
release |
Worker 或应用程序版本 |
将作业参数、凭据、电子邮件地址、文档内容和原始异常文本排除在事件之外。如果操作员需要故障类别,请使用有界 error_type。
仪器结果边界
安装并初始化JavaScript SDK:
npm install telemetry-sh
import telemetry from "telemetry-sh";
telemetry.init("YOUR_API_KEY");
使用相同的时钟测量排队、首次启动和终端完成。成功后发出一行终端行或重试耗尽:
async function processJob(job) {
const startedAtMs = Date.now();
let terminalAttemptStartedAtMs = startedAtMs;
let attemptCount = 0;
try {
const result = await runWithRetry(async () => {
attemptCount += 1;
terminalAttemptStartedAtMs = Date.now();
return performJob(job);
});
await telemetry.log("background_job_completed", {
event_id: crypto.randomUUID(),
job_id: job.id,
job_name: job.name,
queue_name: job.queue,
account_id: job.accountId,
attempt_count: attemptCount,
queue_wait_ms: startedAtMs - job.enqueuedAtMs,
duration_ms: Date.now() - terminalAttemptStartedAtMs,
total_elapsed_ms: Date.now() - startedAtMs,
status: "success",
release: process.env.APP_RELEASE ?? "unknown"
});
return result;
} catch (error) {
await telemetry.log("background_job_completed", {
event_id: crypto.randomUUID(),
job_id: job.id,
job_name: job.name,
queue_name: job.queue,
account_id: job.accountId,
attempt_count: attemptCount,
queue_wait_ms: startedAtMs - job.enqueuedAtMs,
duration_ms: Date.now() - terminalAttemptStartedAtMs,
total_elapsed_ms: Date.now() - startedAtMs,
status: "error",
error_type: classifyJobError(error),
release: process.env.APP_RELEASE ?? "unknown"
});
throw error;
}
}
此示例将终端尝试保留在 duration_ms 中,并将完整重试策略挂起时间保留在 total_elapsed_ms 中。如果您的重试库以不同的方式报告这些边界,请调整计时器,同时保留两个记录的含义。
遥测调用应遵循作业的持久状态更改。决定您的应用程序如何处理遥测传输失败,而不会将已经成功的作业变成其业务副作用的重试。当重新发送相同的终端事件时,在遥测传送重试期间保留 event_id。
队列状态示例
按固定时间间隔从队列的权威状态收集快照:
async function recordQueueSnapshot(queue) {
const state = await queue.inspect();
await telemetry.log("queue_snapshot", {
event_id: crypto.randomUUID(),
queue_name: queue.name,
depth: state.waitingCount,
active_workers: state.activeWorkers,
oldest_wait_ms: state.oldestEnqueuedAtMs
? Date.now() - state.oldestEnqueuedAtMs
: 0,
release: process.env.APP_RELEASE ?? "unknown"
});
}
保持间隔足够频繁,以检测有意义的积压,但不要太频繁,以免相同的样本主导事件量。将零深度记录为零;不要省略它。
将队列等待与运行时分开
此查询将典型等待和尾部等待与尾部执行持续时间进行比较:
SELECT
job_name,
COUNT(*) AS jobs,
approx_percentile_cont(queue_wait_ms, 0.50) AS p50_wait_ms,
approx_percentile_cont(queue_wait_ms, 0.95) AS p95_wait_ms,
approx_percentile_cont(duration_ms, 0.95) AS p95_run_ms
FROM background_job_completed
WHERE timestamp_utc >= now() - INTERVAL '24 hours'
GROUP BY job_name
HAVING COUNT(*) >= 20
ORDER BY p95_wait_ms DESC;
正常运行时间下的高 p95 等待指向容量、调度、优先级或突发处理。运行时间长,正常等待点指向作业代码或依赖项。两者同时上升可能意味着缓慢的工作正在消耗工人的能力并造成积压。
测量重试和终端故障
由于终止事件记录了一个逻辑作业,因此 attempt_count > 1 表示该作业在至少一次重试后恢复:
SELECT
job_name,
COUNT(*) AS jobs,
SUM(CASE WHEN attempt_count > 1 THEN 1 ELSE 0 END) AS retried_jobs,
SUM(CASE WHEN status = 'error' THEN 1 ELSE 0 END) AS permanent_failures,
100.0 * SUM(CASE WHEN attempt_count > 1 THEN 1 ELSE 0 END)
/ NULLIF(COUNT(*), 0) AS retried_job_rate_pct,
100.0 * SUM(CASE WHEN status = 'error' THEN 1 ELSE 0 END)
/ NULLIF(COUNT(*), 0) AS permanent_failure_rate_pct
FROM background_job_completed
WHERE timestamp_utc >= now() - INTERVAL '7 days'
GROUP BY job_name
ORDER BY permanent_failure_rate_pct DESC, retried_job_rate_pct DESC;
如果您每次尝试发出一行,请使用 attempt 字段和稳定的 job_id;分母和解释会有所不同。在每个查询旁边记录行粒度,这样重试就不会意外地被计为单独的客户作业。
检测停滞和丢失的工作
仅终端表无法区分长时间运行的作业和工作线程崩溃后丢失的作业。发出 job_started 加 job_completed 或 job_failed 与相同的 job_id,然后使用左连接查找没有终止事件的开始。
该阈值必须是特定于工作的或源自预期完成窗口。长时间运行的作业可能需要心跳事件。为事件传递延迟添加一个较短的宽限期,以便最新行不会产生误报。
使用完整的 停滞的后台作业查询 而不是将队列深度视为特定作业被卡住的证据。
考虑死信、优先事项和关闭
当重试策略用尽且作业进入死信队列时,记录不同的终端类别或事件。轨迹回放作为新的操作动作与原job_id联动;不要默默地改写原来的失败。
当工作人员不可互换时,按队列或优先级对容量信号进行分段。健康的批量队列可以隐藏总体平均值中阻塞的关键队列。
在部署和关闭期间:
- 在解雇工人之前停止接受新工作;
- 允许记录的换油间隔;
- 区分故意重新排队和执行失败;
- 保留逻辑
job_id并递增尝试状态; - 验证每个启动的作业最终都会产生一个终止事件。
构建仪表板
后台作业仪表板示例 提供 SQL、合成结果和解释。生产环境插桩板应包括:
- 当前队列深度和队列中最早的等待年龄;
- p50 和 p95 按作业名称排队等待;
- p95 按作业名称执行持续时间;
- 重试作业率和永久失败率;
- 当前停滞的作业具有安全相关标识符;
- 按版本划分的数量,以便可以将部署与更改进行比较。
使用完整的时间段来了解趋势。在对费率进行排名之前设置最小交易量规则。仅当工作流程至关重要时,安静队列中的单个失败作业才具有操作重要性;否则,仅按百分比计算,它的排名不应超过大容量回归。
对响应发出告警,而不仅仅是阈值
将每个告警与所有者和操作联系起来:
| 条件 | 可能的问题 | 第一反应 |
|---|---|---|
| 深度和最久等待时间上升 | 需求是否超过容量? | 检查到达率、工人、优先级和依赖性健康状况 |
| 运行时间增加而等待稳定 | 作业代码或依赖项是否会减慢速度? | 比较版本和错误类别 |
| 重试率上升 | 瞬态故障放大有用吗? | 检查有界错误类型和提供者健康状况 |
| 永久性故障增加 | 恢复力竭了吗? | 识别受影响的帐户和死信状态 |
| Start没有终止事件 | 是否有工人坠毁或仪器失踪? | 检查worker、心跳和作业状态 |
避免对一个不完整的存储桶进行分页。需要持续的违规或严重的终端故障,并使阈值与工作流程的预期完成时间保持一致。
生产前验证
运行一个装置:
- 第一次尝试成功;
- 重试成功;
- 永久性故障和死信条目;
- 重复事件交付;
job_started后工人崩溃;- 有心跳的长时间运行的工作;
- 正常耗尽的队列突发;
- 部署关闭并重新排队。
确认逻辑作业计数、尝试计数、队列等待、运行持续时间、丢失的终端事件和受影响的帐户计数。还要验证遥测故障是否无法重播非幂等业务操作。