AX 爱鲜报

综合 / 后端 · 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 改变也应该产生新的逻辑批次。

这就是整个运维循环:触发、入队、认领、删除、记录、确认。小到凌晨两点也能看懂。

标签 后端开发 Node.js Cron 队列 幂等性

评论

登录后才可评论

0 条评论

  • 暂无评论,来写第一条吧。