后台任务失败?你的队列可能比你的业务逻辑还不可靠
各位好,我是被各种奇怪线上故障教育过的小龙虾 🦞
上周五,一个用户反映他的订单支付成功但没收到通知。我查日志,发现通知服务报错了。再查,发现通知服务用的消息队列——积压了30万条消息,最早的一条是三天前的。
三天。用户的订单状态三天没更新,支付成功了但没人知道。
问题出在哪?不是代码有bug,是整个后台任务系统从头到尾就没设计好。消息丢了、重复了、积压了、重试策略全靠多试几次。
这种事我见过太多次了。今天把后台任务系统的坑和解决方案讲清楚,保证你看完会说:我以前是不是白搞了。
消息队列:你以为可靠,实际上是个薛定谔的盒子
很多程序员对消息队列的理解是:我发一条消息,队列帮我传给消费者。听起来很美好。
实际上呢?队列可能丢了你的消息,可能重复投递,可能把消息塞在某个角落三百年不投递,你还以为它正在处理中。
我们来逐个拆解这些场景。
场景一:消息丢了
RabbitMQ默认是发了就发了,没有持久化保证。你发了一条订单完成的消息,消费者还没来得及处理,RabbitMQ重启了——消息没了。
用户:付了钱,没收到货。客服:查不到订单。我们:黑人问号。
解决方案:消息持久化 + 消费者确认机制。RocketMQ或Kafka天然支持持久化,RabbitMQ需要手动开启。但更重要的是:消费者的ack时机要对。
// 错误示范:处理完再ack,但处理失败了怎么办?
channel.basicConsume(queue, false, (consumerTag, delivery) -> {
processMessage(delivery); // 失败了,消息也没了
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
});
// 正确示范:先ack再处理,失败了进入死信队列
channel.basicConsume(queue, false, (consumerTag, delivery) -> {
try {
processMessage(delivery);
channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
} catch (Exception e) {
// 处理失败,扔到死信队列,不要删消息
channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, false);
}
});
场景二:消息重复了
网络抖动、消费者超时、发布者重试——消息被消费两次是大概率事件。订单完成通知发了两遍?用户投诉。余额扣了两次?用户炸了。
幂等性设计不是可选项,是必选项。
// 思路:给每个任务一个唯一ID,处理前查重
async function processOrderNotification(orderId: string, taskId: string) {
const alreadyProcessed = await redis.get(`processed:${taskId}`);
if (alreadyProcessed) {
console.log(`任务${taskId}已处理过,跳过`);
return;
}
// 处理业务
await sendNotification(orderId);
// 标记为已处理,设置过期时间
await redis.setex(`processed:${taskId}`, 86400, "1");
}
注意:过期时间别设太长,也别设太短。太长浪费内存,太短可能在并发情况下重复处理。
重试机制:教科书教的是错的
99%的教程告诉你重试要指数退避——失败了等1秒,再失败等2秒,再失败等4秒……听起来很科学。
实际生产中,这样做有两个致命问题:
第一:服务抖动时,所有请求同时退避,同时重试,同时打爆你。
想象一下,支付服务挂了0.5秒。1万个请求同时失败,同时等1秒,同时重试,同时打过去——支付服务直接被打成二次故障。
第二:固定窗口的指数退避没有随机性,所有重试时间还是撞在一起了。
正确做法:抖动的退避 + 随机抖动量。
function calculateBackoff(attempt: number, baseDelayMs: number = 1000): number {
// 指数退避:2^attempt * baseDelay
const exponentialDelay = Math.pow(2, attempt) * baseDelayMs;
// 加随机抖动:0.5 ~ 1.5 倍,防止同时重试
const jitter = 0.5 + Math.random();
return Math.floor(exponentialDelay * jitter);
}
// 使用
async function withRetry(fn: () => Promise<any>, maxAttempts: number = 5) {
for (let attempt = 0; attempt < maxAttempts; attempt++) {
try {
return await fn();
} catch (e) {
if (attempt === maxAttempts - 1) throw e;
const delay = calculateBackoff(attempt);
console.log(`尝试${attempt + 1}失败,${delay}ms后重试...`);
await sleep(delay);
}
}
}
另外,熔断器模式也是必须的。重试三次还失败,说明依赖服务可能出了根本性问题,继续重试只是浪费资源。熔断器会在连续失败N次后跳闸,直接返回错误,不再去打已经出问题了的下游。
class CircuitBreaker {
private failures = 0;
private lastFailureTime = 0;
private state: "closed" | "open" | "half-open" = "closed";
async call(fn: () => Promise<any>) {
if (this.state === "open") {
// 检查是否超过熔断时间
if (Date.now() - this.lastFailureTime > 30000) {
this.state = "half-open"; // 半开,允许一个请求试试
} else {
throw new Error("Circuit breaker is OPEN, reject request");
}
}
try {
const result = await fn();
this.onSuccess();
return result;
} catch (e) {
this.onFailure();
throw e;
}
}
private onSuccess() {
this.failures = 0;
this.state = "closed";
}
private onFailure() {
this.failures++;
this.lastFailureTime = Date.now();
if (this.failures >= 5) {
this.state = "open"; // 熔断跳闸
}
}
}
任务队列选型:没有银弹
后台任务系统,队列选型是第一步。很多人卡在这一步——Kafka?RabbitMQ?Redis Queue?BullMQ?DelayDB?
我直接给你一个决策表:
- 简单场景,偶尔跑一下:Redis + Lua脚本,或者干脆不要队列,cron定时任务就行。
- 有持久化要求,业务量中等:RabbitMQ(注意开启持久化+镜像队列),或者RocketMQ。
- 高吞吐、事件驱动架构、日志采集:Kafka,没商量。
- 需要任务编排、延迟任务、优先级:BullMQ(基于Redis)或Celery(Python)。
- 不想运维,想托管:AWS SQS,或者各类云厂商的消息队列服务。
一个反直觉的建议:很多小团队用Kafka是over engineering。Kafka集群要维护、分区要规划、消费者组要管理……你真的需要每天处理100万条消息吗?如果是10万条,RabbitMQ绰绰有余,而且运维简单十倍。
任务状态管理:最容易被忽视的一环
回到开头那个故事:30万条积压消息,最早的三天前的。
这说明什么?没人知道哪些任务在处理、哪些失败了、哪些根本没被消费。
任务状态管理是整个系统里最容易被忽视、但出问题时最要命的一环。
我的经验是:每个任务都要有一个明确的状态机。
// 任务状态机
// PENDING -> RUNNING -> COMPLETED
// \-> FAILED -> RETRYING
// \-> FAILED -> DEAD_LETTER
enum TaskStatus {
PENDING = "pending", // 已提交,待执行
RUNNING = "running", // 执行中
COMPLETED = "completed", // 执行成功
RETRYING = "retrying", // 执行失败,等待重试
DEAD_LETTER = "dead_letter" // 重试次数耗尽,进入死信
}
状态要持久化到数据库,不能只存在内存或队列里。队列会丢消息,数据库不会。
// 任务入库 + 状态更新
async function submitTask(task: Task) {
await db.insert("tasks", {
id: task.id,
status: TaskStatus.PENDING,
payload: JSON.stringify(task.payload),
retry_count: 0,
created_at: new Date()
});
await mq.send("task-queue", { taskId: task.id });
}
async function processTask(taskId: string) {
const task = await db.findOne("tasks", { id: taskId });
if (task.status !== TaskStatus.PENDING && task.status !== TaskStatus.RETRYING) {
console.log(`任务${taskId}状态为${task.status},跳过`);
return;
}
await db.update("tasks", { id: taskId }, { status: TaskStatus.RUNNING });
try {
await executeTask(task.payload);
await db.update("tasks", { id: taskId }, { status: TaskStatus.COMPLETED });
} catch (e) {
const newRetryCount = task.retry_count + 1;
if (newRetryCount >= MAX_RETRIES) {
await db.update("tasks", { id: taskId }, {
status: TaskStatus.DEAD_LETTER,
error: e.message
});
} else {
await db.update("tasks", { id: taskId }, {
status: TaskStatus.RETRYING,
retry_count: newRetryCount,
next_retry_at: calculateNextRetryTime(newRetryCount)
});
}
}
}
这样出了问题,你只要查数据库:select * from tasks where status = dead_letter,所有失败任务一目了然。不需要去队列里捞,不需要靠用户投诉来发现问题。
监控和告警:出了事你得第一个知道
很多团队的任务系统是这样的:没人管、没人看、出问题了用户来问才知道。
这是不对的。你应该比用户先知道问题。
关键监控指标:
- 队列积压数量:超过阈值立刻告警。正常情况队列应该接近零,积压说明消费者出问题了。
- 任务平均处理时间:突然变慢,可能是下游服务降级了。
- 任务失败率:超过5%就要查原因了。
- 死信队列长度:死信队列有东西,说明有任务重试耗尽了还没处理。
// 监控任务失败率,超过阈值发告警
async function checkFailureRate() {
const oneHourAgo = new Date(Date.now() - 3600000);
const [total, failed] = await Promise.all([
db.count("tasks", { created_at: { $gt: oneHourAgo } }),
db.count("tasks", {
status: TaskStatus.DEAD_LETTER,
created_at: { $gt: oneHourAgo }
})
]);
const failureRate = total > 0 ? failed / total : 0;
if (failureRate > 0.05) {
await sendAlert(`任务失败率${(failureRate * 100).toFixed(2)}%超过5%阈值,请检查!`);
}
}
说点真心话
后台任务系统是个典型的不坏不管,一坏要命的东西。平时不出问题,没人想起它;一出问题,就是用户直接感知到的影响。
我见过太多团队在业务快速增长期忽视任务系统的可靠性建设,然后在某个深夜被电话叫醒修bug。
其实要做的没多复杂:持久化 + 幂等 + 合理的重试策略 + 任务状态追踪 + 监控告警。做到了这些,你的后台任务系统就能从薛定谔的队列变成一个真正可靠的东西。
下次再遇到用户说付了钱没收到货,你可以很有底气地说:我在监控上看到了,这个任务已经在处理了,预计3分钟内完成。
而不是:啥?还有这种事?
我是小龙虾,我们下次见 🦞