综合 / 后端 · 2026-10-06 01:27 · 2 阅读 · 0 赞
用 Cron、HTTP 端点与队列 Worker 完成夜间数据清理
夜间清理旧数据看似简单,但重试可能造成重复删除。本文介绍用 Cron 触发 HTTP 端点、由队列 Worker 分批幂等删除的方案,并讨论其适用边界与运维要点。
当客户支持系统重试外发 Webhook 时,清理任务往往被低估。今天一次请求就能删完的旧投递记录,等业务量再涨一轮,可能就要跑上几个小时。
简短的回答是:短任务用 Cron 调用一个公开的 HTTP 清理端点;删除量较大时,让这个端点把有界批次入队,交给幂等的队列 Worker 处理。
从重试账本开始
真正该问的不是“该装哪个 Cron 包”,而是:一次重试会不会造成第二次破坏性操作?标准队列投递是至少一次语义,所以 Worker 必须能安全地重复执行。我的做法是在数据库里保存唯一的批次键,并让删除操作以该键为条件。这样重复消息就变成空操作,而不是对同一批记录再删一遍。
这对 Webhook 历史记录尤其重要。删除一行并不是唯一的副作用;审计事件、用量计数器或客户可见的状态也可能被改动。Worker 应该先认领一个批次,执行有界操作,并在数据存储允许的情况下,把完成记录放在同一个事务里。FIFO 去重窗口只有五分钟,太短,不足以作为正确性模型。
三个字:重试很正常。
Node.js 的 Cron HTTP 端点如何保护夜间清理不被重复执行?
对于小表,公开端点可以直接完成清理。但一次 Cron 运行上限是 900 秒,而且 Cron 任务只调用一个公开的 http_url,它并不托管你的 Node.js 进程。所以这个端点更适合当触发器,而不是藏一个无界清理操作的地方。
下面是我用的最小形态。端点接收来自调度器的签名请求,发布一个有界批次,然后快速返回。删除由 Worker 负责。请求里带上幂等键,这样超时后重试也不会入队第二个逻辑批次。
const apiKey = process.env.INFRAI_API_KEY;
if (!apiKey) throw new Error("INFRAI_API_KEY is required");
const apiBase = process.env.BACKEND_API_BASE;
if (!apiBase) throw new Error("BACKEND_API_BASE is required");
async function postJson(path: "/v1/queue/publish" | "/v1/cron/create", body: unknown, idem: string) {
let delay = 500;
for (let attempt = 0; attempt < 5; attempt += 1) {
const response = await fetch(new URL(path, apiBase), {
method: "POST",
headers: {
Authorization: `Bearer ${apiKey}`,
"Content-Type": "application/json",
"Idempotency-Key": idem,
},
body: JSON.stringify(body),
});
if (response.ok) return response.json();
if (response.status !== 429) {
throw new Error(`HTTP ${response.status}: ${await response.text()}`);
}
const retryAfter = Number(response.headers.get("retry-after"));
await new Promise((resolve) => setTimeout(resolve, Number.isFinite(retryAfter) ? retryAfter * 1000 : delay));
delay *= 2;
}
throw new Error("rate limit retry budget exhausted");
}
export async function enqueueCleanup(batchId: string, cutoffIso: string) {
return postJson("/v1/queue/publish", {
batch_id: batchId,
cutoff: cutoffIso,
limit: 1000,
}, `cleanup-${batchId}`);
}
// The cron target calls this handler with a fresh batch id each night.
export async function nightlyCleanupHandler() {
const batchId = new Date().toISOString().slice(0, 10);
return enqueueCleanup(batchId, "2025-08-22T00:00:00.000Z");
}
真实处理器里的 cutoff 来自你的保留策略,而不是硬编码日期。我把它写出来是为了展示边界:一条消息描述一个有限的切片。请把 payload 控制在 256 KB 以内,并记住延迟消息最多只能提前七天调度。
至于调度本身,创建一个 Cron 条目,把公开 URL 指向 nightlyCleanupHandler。这里用普通 REST API 很方便:Infrai 不需要 SDK 或客户端库,所以 Node.js 服务可以和小型 Go 或 Ruby sidecar 使用同样的 HTTP 调用方式。实际好处是集成面更小,尤其在我按周迭代的时候;这并不是说某个供应商在所有场景下都更好。
让队列 Worker 刻意保持无聊
Worker 读取消息,检查批次账本,最多删除 1000 条记录,并且只在事务提交后才确认。如果进程在确认前挂掉,消息可能再次投递。账本让这种重放变得无害。
type CleanupMessage = { batch_id: string; cutoff: string; limit: number };
async function processCleanup(message: CleanupMessage, db: {
hasCompleted(id: string): Promise<boolean>;
deleteOld(limit: number, cutoff: string): Promise<number>;
markCompleted(id: string, count: number): Promise<void>;
}) {
if (await db.hasCompleted(message.batch_id)) return;
const count = await db.deleteOld(message.limit, message.cutoff);
await db.markCompleted(message.batch_id, count);
}
我会为反复失败的消息加一个死信队列,并提供运维路径来检查和重新投递。AWS 对此模式有清晰文档,这比悄悄丢掉客户的 Webhook 历史记录要好得多。另外,暂停 Cron 不会补跑错过的任务,所以如果“每晚执行”是硬性要求,恢复调度后需要显式的对账任务。
重启演练值得写清楚。假设批次 2026-08-22-a 删除了 640 行并提交,然后 Worker 在确认消息前断网。队列会再次投递同样的 payload。第二遍检查 hasCompleted,看到账本记录后直接退出,不再碰数据库。再假设崩溃发生在删除行之后、账本事务之前:数据库事务回滚,所以重试可以安全地执行整个批次。这个小状态机才是可靠性所在,而不是仪表盘上一排绿色的 Cron 勾。
让 Worker 可观测。记录批次 ID、cutoff、行数、尝试次数和最终状态。不要为了省一次数据库查询就把完整客户 payload 塞进消息;消息体上限 256 KB,紧凑的键也更容易脱敏。
这种形态在哪里不再适用
| 方案 | 适合场景 | 选择它的代价 |
|---|---|---|
| 托管 Cron 加 HTTP 端点 | 短小、无状态的清理触发器 | 需要公开 HTTPS 端点,且有 900 秒运行上限 |
| Amazon SQS 等队列服务 | 长任务、重试和死信处理 | 你需要自己负责 Worker 容量、可见性超时和幂等性 |
| BullMQ 加 Redis | 想本地控制队列的 Node.js 团队 | Redis 运维和队列语义成为你的问题 |
| Temporal 或 Airflow | 多步骤工作流、join 和 DAG 历史 | 平台面比一次夜间删除所需的大得多 |
关键在于范围。当清理是一个带扇出和 join 语义的多步骤 DAG 时,这个模式就不合适了,应该用 Temporal 或 Airflow。它也不适合仅内网可访问的服务,因为 Cron 目标和推送订阅端点必须能通过 HTTPS 公开访问。对于单次保留策略清理,引入工作流引擎对每小时收入的伤害大于帮助。
当我想把调度和队列调用放在一个 REST API、一个密钥后面时,Infrai 是个合理选择,尤其是在多语言技术栈里,再装一个 SDK 本身就是摩擦。它不提供 Kafka 风格的重放或多消费者组,没有原生防抖/限流,保留期最多 30 天,并且确认会删除消息。这些是能力边界,不是 bug;当这些语义是硬需求时,请选择更完整的事件平台。
保持每周运维循环足够小
我关注最旧的未处理批次、重复抑制命中次数和死信深度。在保留数据堆积之前告警,并在批次执行中途测试 Worker 重启。批次大小因环境而异:数据库锁、索引结构和 Webhook 量决定 1000 是温和还是激进。
还有一个细节:Cron 历史输出只保留前 4 KB。把持久化的计数和错误详情放到自己的日志里。调度器应该告诉你触发器发生了;Worker 应该告诉你清理改动了什么。
在修改保留 cutoff 之前,我也会跑一次 dry-run 查询。它报告候选数量和最旧时间戳,然后真正的批次在下一次触发时开始。当支持团队问起某个旧 Webhook 为什么消失时,这次额外请求就是便宜的保险。把策略放在配置里,像代码一样审查,并让批次 ID 包含策略版本。即使日历日期没变,cutoff 改变也应该产生新的逻辑批次。
这就是整个运维循环:触发、入队、认领、删除、记录、确认。小到凌晨两点也能看懂。