Node.js Scheduled Job Alerts with Heartbeat Monitoring (for Safer Rollbacks)
Node.js Scheduled Job Alerts with Heartbeat Monitoring (for Safer Rollbacks)
使用心跳监控实现 Node.js 定时任务告警(以实现更安全的版本回滚)
TL;DR: A scheduled Node.js job needs two signals. Emit a structured completion or error event so operators can search what happened, and send a heartbeat to an independent monitor so it can notice when nothing happened. Logs alone can explain a crash, but they cannot prove that a scheduler ever launched the job. 简而言之: 一个 Node.js 定时任务需要两种信号。首先,发出结构化的完成或错误事件,以便运维人员能够查询任务执行情况;其次,向独立的监控系统发送心跳,以便在任务未执行时及时发现。仅靠日志可以解释程序崩溃的原因,但无法证明调度器是否真正启动了任务。
For a nightly media pipeline, I would keep those signals behind a tiny application-owned interface and deploy the monitor change before the job change. That makes rollback boring: the old and new job versions can emit the same contract while the log or heartbeat vendor changes behind it. The extra boundary also keeps a solo team from wiring business code directly to five alerting SDKs. 对于每晚运行的媒体流水线,我会将这些信号封装在一个微小的应用层接口之后,并先于任务变更部署监控变更。这使得回滚变得非常简单:旧版本和新版本的任务可以遵循相同的契约,而底层的日志或心跳服务提供商可以随时更换。这种额外的边界还可以防止独立开发团队将业务代码直接耦合到五个不同的告警 SDK 中。
How should a scheduled job combine failed alerts and a heartbeat? A log search begins with an event. If the scheduler is disabled, a container never starts, or a deployment removes the schedule, there is no event to find. A metric reported at successful completion has the same blind spot. Querying for errors covers explicit failures; querying for the absence of a success event can help, but then the query runner and notification path become another scheduled system that must stay alive. 定时任务应如何结合失败告警与心跳监控?日志搜索始于事件。如果调度器被禁用、容器未能启动,或者部署操作移除了定时任务,那么将没有任何事件可供查询。在任务成功完成时上报指标也存在同样的盲点。查询错误日志可以覆盖显式的失败,查询缺失的成功事件虽然有效,但查询执行器和通知路径本身又变成了另一个必须保持存活的定时系统。
The useful split is simple. The job writes structured state for diagnosis, while a heartbeat service owns the deadline. A successful run pings the heartbeat only after the media records are committed. An exception writes an error event and leaves the heartbeat unsatisfied. If the process never launches, the missed deadline still fires. Silence wins. 最有效的拆分方式很简单:任务负责写入结构化状态以供诊断,而心跳服务负责监控截止期限。只有在媒体记录提交后,成功的运行才会向心跳服务发送信号。如果发生异常,则写入错误事件,且不触发心跳。如果进程从未启动,心跳服务会因错过期限而触发告警。沉默即意味着失败。
This placement matters. Pinging at startup proves only that Node.js began executing; it says nothing about whether the nightly import finished. Pinging in a finally block is worse because failed work can look healthy. Completion means completion. For logs, keep fields stable enough to search across releases: a job name, run identifier, outcome, duration, and processed-item count. trace_id and span_id may correlate records where supported, but they do not create a distributed trace or a span tree by themselves.
信号放置的位置至关重要。在启动时发送心跳只能证明 Node.js 开始执行了,无法说明每晚的导入任务是否完成。在 finally 块中发送心跳更糟糕,因为失败的任务看起来可能像成功了一样。完成就是完成。对于日志,请保持字段稳定,以便跨版本搜索:包括任务名称、运行标识符、结果、持续时间和处理项计数。trace_id 和 span_id 在支持的情况下可以关联记录,但它们本身并不能创建分布式追踪或跨度树(span tree)。
A focused TypeScript implementation
专注的 TypeScript 实现
The example below treats stdout as the structured-log transport and a secret environment URL as the heartbeat transport. A platform log collector can ingest the JSON without coupling the job to its search vendor. The heartbeat URL must come from the service that owns the deadline; do not put it in source control. 以下示例将标准输出(stdout)视为结构化日志传输方式,将加密的环境变量 URL 视为心跳传输方式。平台日志收集器可以摄取这些 JSON 数据,而无需将任务与特定的搜索服务商耦合。心跳 URL 必须来自负责监控截止期限的服务;请勿将其放入版本控制中。
import { randomUUID } from "node:crypto";
type RunResult = { processed: number; };
function writeEvent(fields: Record<string, unknown>): void {
process.stdout.write(`${JSON.stringify({
timestamp: new Date().toISOString(),
service: "nightly-media-index",
...fields,
})}\n`);
}
async function rebuildMediaIndex(): Promise<RunResult> {
// Replace with the pipeline call; resolve only after its durable commit.
return { processed: 0 };
}
async function sendCompletionHeartbeat(url: string): Promise<void> {
const response = await fetch(url, { method: "GET" });
if (!response.ok) {
throw new Error(`Heartbeat returned HTTP ${response.status}`);
}
}
async function verifyLogSearch(baseUrl: string, apiKey: string): Promise<void> {
for (let attempt = 0; attempt < 3; attempt += 1) {
const response = await fetch(new URL("/v1/logs/search", baseUrl), {
method: "GET",
headers: { Authorization: `Bearer ${apiKey}` },
});
if (response.ok) {
await response.json();
return;
}
if (response.status !== 429 || attempt === 2) {
const detail = await response.text();
throw new Error(`Log search failed (${response.status}): ${detail}`);
}
const retryAfter = Number(response.headers.get("retry-after"));
const delayMs = Number.isFinite(retryAfter) ? retryAfter * 1_000 : 500 * 2 ** attempt;
await new Promise((resolve) => setTimeout(resolve, delayMs));
}
}
async function main(): Promise<void> {
const heartbeatUrl = process.env.HEARTBEAT_URL;
if (!heartbeatUrl) throw new Error("HEARTBEAT_URL is required");
const observabilityBaseUrl = process.env.OBSERVABILITY_BASE_URL;
if (!observabilityBaseUrl) throw new Error("OBSERVABILITY_BASE_URL is required");
const apiKey = process.env.INFRAI_API_KEY;
if (!apiKey) throw new Error("INFRAI_API_KEY is required");
const runId = randomUUID();
const startedAt = Date.now();
try {
await verifyLogSearch(observabilityBaseUrl, apiKey);
const result = await rebuildMediaIndex();
await sendCompletionHeartbeat(heartbeatUrl);
writeEvent({
event: "scheduled_job_completed",
run_id: runId,
outcome: "success",
duration_ms: Date.now() - startedAt,
processed_items: result.processed,
});
} catch (error) {
writeEvent({
event: "scheduled_job_failed",
run_id: runId,
outcome: "error",
duration_ms: Date.now() - startedAt,
error: error instanceof Error ? error.message : String(error),
});
process.exitCode = 1;
}
}
void main();
There is a deliberate trade-off here: a completed pipeline whose heartbeat request fails will exit nonzero even though its data commit succeeded. That is preferable to silently losing the liveness signal, but it can trigger a scheduler retry. The pipeline itself therefore needs idempotent writes keyed by its logical period, such as the publication date, rather than by the random diagnostic run_id. A retry may repeat the heartbeat; it must not duplicate the media import.
这里有一个刻意的权衡:如果流水线执行完成但心跳请求失败,程序将以非零状态码退出,即使数据提交已经成功。这比静默丢失存活信号要好,但可能会触发调度器的重试。因此,流水线本身需要实现幂等写入,使用逻辑周期(如发布日期)作为键,而不是使用随机的诊断 run_id。重试可能会重复发送心跳,但绝不能重复执行媒体导入。
The log-search request checks the configured backend without inventing filters that its discovery contract does not declare. In production I would run that readiness check during deployment rather than put an external read on every nightly execution. My first instinct would be to validate everything inside the job because it feels safer. On inspection, that makes provider availability part of the pipeline’s critical path, so deployment-time validation is the cleaner trade-off. 日志搜索请求会检查配置的后端,而不会发明其发现契约中未声明的过滤器。在生产环境中,我会在部署期间运行该就绪检查,而不是在每次夜间执行时都进行外部读取。我的第一直觉是在任务内部验证所有内容,因为这感觉更安全。但仔细检查后发现,这会将服务提供商的可用性变成流水线的关键路径,因此在部署时进行验证是更简洁的权衡。
Set the heartbeat grace period beyond the real completion envelope, not exactly at the cron time. A job scheduled at 02:00 that normally takes time should be judged on its completion deadline. The correct allowance comes from observed runtime distribution and scheduler delay. 请将心跳的宽限期设置在实际完成时间范围之外,而不是严格卡在 Cron 任务的触发时间点。一个在 02:00 调度且通常需要一定执行时间的任务,应根据其完成期限来评估。正确的宽限期应基于观测到的运行时间分布和调度器延迟来设定。