RabbitMQ 异步调度落地路径(评测批量任务的队列方案)

AI 大模型评测平台一次评测会拆成成百上千条子任务并发打给不同模型,若在 Web 请求线程内同步等待,连接会超时、线程被长期占用。把评测任务投到 RabbitMQ 队列,由独立消费者按自身节奏拉取执行,才能做到请求秒回、任务削峰、失败重试与多实例负载分发。下文对比线程池与消息队列的差异,给出 Spring Boot 可运行的队列与消费者代码。

说明:本文基于通用工程实践对平台调度层做合理推演,非平台真实源码。方案适用于任意 Spring Boot + RabbitMQ 技术栈,读者可按实际模型供应商与限流策略调整路由与并发参数。

一、为什么评测调度要离开请求线程

批量评测的请求特征是”提交快、执行慢”。一道综合评测题可能要调 5 家模型、每家 3 次重复,单次 HTTP 往返几百毫秒到数秒。同步处理会让 Tomcat 线程池被评测任务占满,正常的管理接口也拿不到线程。

把任务提交和任务执行分离:接口只做”接收题目、生成任务、入队”三件事,立刻返回任务 ID;真正的模型调用交给后台消费者。这种架构让 Web 层保持轻量,评测压力被队列吸收。

二、消息队列与线程池的核心差异

二者都能做异步,但定位不同。线程池是进程内的任务调度工具,消息队列是跨进程的通信与解耦工具。

对比维度 线程池(ThreadPoolExecutor) RabbitMQ 消息队列
可靠性 JVM 重启未跑任务丢失 消息持久化,Broker 重启不丢
分布式 仅本机,无法跨实例分发 多消费者实例共同分担
削峰 队列满则阻塞或拒绝 任务堆积,消费者按速处理
重试 需自己写补偿逻辑 消费失败可重新入队(Nack)
可观测 无天然积压视图 管理台看得到队列深度
适用 单实例、可丢、低延迟内部任务 不能丢、需跨实例、批量评测

评测平台多实例部署且任务不允许静默丢失,消息队列是更稳的选择。线程池适合”写访问日志””刷新本地缓存”这类丢了不致命的内部杂活。

三、RabbitMQ 队列方案设计

3.1 交换机与队列拓扑

评测任务按”优先级”或”模型供应商”路由到不同队列。普通批量任务进 eval.normal,重要回归测试进 eval.priority,用 basic.qos(prefetch=1) 保证公平分发,避免某个消费者被长任务占死。

3.2 任务投递与持久化

投递时声明队列 durable、消息 delivery_mode=2,确保 RabbitMQ 重启后任务还在。

@Configuration
public class RabbitConfig {
    @Bean
    public Queue evalNormal() {
        return new Queue("eval.normal", true); // durable=true
    }

    @Bean
    public DirectExchange evalExchange() {
        return new DirectExchange("eval.exchange", true, false);
    }

    @Bean
    public Binding bindNormal(Queue evalNormal, DirectExchange evalExchange) {
        return BindingBuilder.bind(evalNormal).to(evalExchange).with("normal");
    }
}
@Service
public class EvalProducer {
    private final RabbitTemplate rabbit;

    public EvalProducer(RabbitTemplate rabbit) { this.rabbit = rabbit; }

    public void submit(EvalTask task) {
        rabbit.convertAndSend("eval.exchange", "normal", task,
            message -> {
                message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
                return message;
            });
    }
}

3.3 消费者与手动 Ack

消费者拉到消息后执行业务,成功才 basicAck;执行中崩溃则 RabbitMQ 自动把消息重新投递给其他实例,天然支持重试。

@Component
public class EvalConsumer {
    @RabbitListener(queues = "eval.normal", ackMode = "MANUAL")
    public void handle(EvalTask task, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException {
        try {
            runEvaluation(task);          // 调多家模型、比对结果
            channel.basicAck(tag, false); // 成功确认,消息出队
        } catch (TransientException e) {
            channel.basicNack(tag, false, true); // 可重试,重回队列
        } catch (FatalException e) {
            channel.basicNack(tag, false, false); // 不可重试,进入死信
        }
    }
}

四、从提交到完成的端到端流程

  1. 管理端提交评测集,平台按题×模型展开成子任务列表;
  2. 每个子任务序列化为 EvalTask 投到 eval.exchange;
  3. 多个消费者实例竞争拉取,prefetch=1 保证匀速;
  4. 消费者调用模型 API,把结果写回数据库,更新任务进度;
  5. 全部子任务完成,聚合生成评测报告并通知提交人。

这套链路的关键收益有三处:突发提交时队列缓冲,模型 API 限流也不至于冲垮平台;某消费者进程崩溃,未 Ack 的任务自动漂移到健康实例;横向扩消费者即可线性提升评测吞吐。

五、踩坑点

自动 Ack 是头号陷阱:RabbitMQ 把消息发出即视为成功,消费者中途崩溃任务就丢了,必须手动 Ack 且放在业务逻辑末尾。队列深度要监控,Ready 数持续上涨说明消费者算力不足,应及时扩容实例。评测任务要设计幂等键(题目 ID + 模型 + 轮次),防止重复消费导致重复计费。

六、死信队列与延迟重试

瞬时故障(模型 API 限流、网络抖动)不应立即判死刑。把可重试异常 Nack 且 requeue=true 会让消息立刻重回队首反复重试,容易打满队列。更稳的做法是投到死信交换机,配一个带 TTL 的延迟队列,让消息”睡几秒”再回来,实现指数退避。

失败类型 处理 去向
限流 / 超时 延迟重试 死信队列 + TTL 重回
参数错误 标记失败 死信,不重试
模型返回异常 有限重试 重试 3 次后死信

延迟队列依赖 RabbitMQ 的死信机制:主队列设置 x-dead-letter-exchange 与 x-message-ttl,消息过期后自动转到重试队列。配合消费端的重试计数,超过上限就入永久死信供人工排查,既不丢任务也不阻塞主链路。

七、与成本监控的衔接

每个子任务消费时调用模型 API,返回里带 usage 的 token 明细。消费者在写评测结果的同时,把消耗上报给成本计量模块累加(见 Redis 实时成本监控一文)。这样”调度—执行—计费”形成一条链:队列管执行节奏,Redis 管花了多少钱,预算超限时反向通知调度层暂停新任务入队,构成闭环控制。

常见问题(FAQ)

Q1:线程池比 RabbitMQ 简单,能用吗?

单实例、任务可丢时用线程池;多实例、任务不可丢必须上队列。

Q2:消费者崩溃任务会丢吗?

用手动 Ack,崩溃未确认的消息会被重新投递给其他实例。

Q3:怎么防止重复评测计费?

给子任务加幂等键(题 ID+模型+轮次),消费前先查已处理标记。

版权声明:本文内容由互联网用户自发贡献,该文观点仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 qiqicto@qq.com 举报,一经查实,本站将立刻删除。
赞 (0)
套餐观察官的头像套餐观察官普通用户

相关推荐

返回顶部