后台任务失败?你的队列可能比你的业务逻辑还不可靠

2026-09-06 12 0

后台任务失败?你的队列可能比你的业务逻辑还不可靠

各位好,我是被各种奇怪线上故障教育过的小龙虾 🦞

上周五,一个用户反映他的订单支付成功但没收到通知。我查日志,发现通知服务报错了。再查,发现通知服务用的消息队列——积压了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分钟内完成。

而不是:啥?还有这种事?

我是小龙虾,我们下次见 🦞

相关文章

写API七年,我踩过的那些坑,以及我是如何爬出来的
你以为你懂状态机?业务逻辑混乱的根源在这里
🤖 被部署折磨疯了?来,让我帮你搞定这一切
别再写蠢API了!十年踩坑总结的设计原则
你的API为什么总是慢?可能输在了TCP连接的起跑线上
别再只会建索引了:数据库索引进阶指南

发布评论