大模型网关削峰填谷:基于消息队列的异步任务池
在大模型系统落地到实际生产业务的过程中,除了智能对话、搜索补全这类必须在数百毫秒内返回的首字交互场景外,有大量业务天然属于长耗时、高计算密度的离线或半离线任务。典型的场景包括:企业级合同长文档合规审查、全量代码库安全扫描与重构建议、数十万字营销文案的批量生成,以及企业私域知识库的大规模切片与向量化嵌入。
这类任务在单次推理或多次链式调用(Agent Loop)中,耗时往往从数十秒延伸至数分钟。如果网关层依然沿用传统的同步 HTTP 或 RPC 调用模式,客户端保持长连接阻塞等待,系统的网络文件描述符(FD)、Tomcat 工作线程池以及微服务连接池会在短时间内被耗尽。更致命的是,一旦上游业务在某一时刻批量提交成百上千个文档,突发的并发流量会瞬间击穿网关,向上游大模型服务商发起雪崩式请求,触发大面积的 HTTP 429 Too Many Requests 错误,导致整体业务瘫痪。
解决这一工程痛点的标准解法,是在大模型网关中引入基于消息队列(Message Queue)的异步任务池架构,通过“异步接收、状态落库、消息缓冲、受控消费、结果通知”的全流程闭环,实现流量的高效削峰填谷。
架构设计:从同步阻塞到异步流水线
在异步削峰架构中,核心思想是将客户端的“任务提交请求”与底层的“推理计算执行”彻底解耦。整个架构由四个核心组件构成:
- 接入网关(API Gateway):负责鉴权、参数校验、生成全局唯一任务 ID、初始化任务状态机,并将任务载荷投递到消息队列,随后立即向客户端返回 HTTP 202 Accepted 状态码及任务凭证。
- 消息中枢(RocketMQ / Kafka):作为削峰填谷的蓄水池,承载瞬时峰值流量,隔离前后端处理速度的不对称性,提供可靠的消息持久化和顺序保证。
- AI Worker 计算集群:作为消息消费者,根据预设的 RPM(Requests Per Minute)和 TPM(Tokens Per Minute)配额,以平滑可控的速率拉取任务,调用大模型接口并完成结果后处理。
- 状态与通知服务:基于 Redis 和关系型数据库维护任务状态机,支持客户端主动轮询、长轮询、SSE/WebSocket 实时推送以及 Webhook 异步回调。
+-----------------------------------------------------------------------------------+ | 基于 RocketMQ 的大模型异步削峰架构 | +-----------------------------------------------------------------------------------+ [客户端 / 上游业务系统] | 1. POST /api/v1/ai/tasks (提交长文本审查/批量生成) v [大模型接入网关] ---> 2. 状态机初始化 (Redis / MySQL 记录 PENDING) | 3. 投递任务载荷到 RocketMQ | 4. 立即返回 TaskId (HTTP 202 Accepted, 耗时 < 15ms) v [RocketMQ 任务中枢 (Topic: LLM_TASK_DISPATCH_TOPIC)] | | 5. Worker 按令牌桶/配额平滑拉取消息 (Rate-Controlled Pulling) v [AI Worker 消费集群] <---> [上游大模型 API (受控 RPM/TPM,严格防 429)] | | 6. 推理完成,更新状态机为 SUCCESS,写入结果载荷 v [Redis / 数据库持久化] ---> 7. 通过 Webhook 回调 / SSE 推动结果给客户端接入层实现:快速握手与状态持久化
接入层的关键在于“轻量与极速”。网关节点不承担任何重型计算,单次请求的处理耗时必须压缩在 15 毫秒以内。请求到达后,生成雪花算法 ID 或 UUID,将初始状态与元数据写入 Redis 哈希结构,投递消息后直接响应。
package com.example.gateway.controller; import com.example.gateway.domain.AiTaskRequest; import com.example.gateway.domain.TaskResponse; import com.example.gateway.domain.TaskStatus; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.apache.rocketmq.client.producer.SendCallback; import org.apache.rocketmq.client.producer.SendResult; import org.apache.rocketmq.spring.core.RocketMQTemplate; import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.http.HttpStatus; import org.springframework.http.ResponseEntity; import org.springframework.messaging.support.MessageBuilder; import org.springframework.web.bind.annotation.*; import java.time.Duration; import java.util.HashMap; import java.util.Map; import java.util.UUID; @Slf4j @RestController @RequestMapping("/api/v1/ai/tasks") @RequiredArgsConstructor public class AiTaskController { private final RocketMQTemplate rocketMQTemplate; private final StringRedisTemplate redisTemplate; private static final String TOPIC_LLM_TASK = "LLM_TASK_DISPATCH_TOPIC"; private static final Duration TASK_TTL = Duration.ofDays(3); @PostMapping("/async-submit") public ResponseEntity<TaskResponse> submitTask(@RequestBody AiTaskRequest request) { String taskId = "TASK_" + UUID.randomUUID().toString().replace("-", ""); String taskKey = "llm:task:" + taskId; // 1. 初始化任务状态机,写入 Redis Map<String, String> meta = new HashMap<>(); meta.put("status", TaskStatus.PENDING.name()); meta.put("userId", request.getUserId()); meta.put("bizType", request.getBizType()); meta.put("createdAt", String.valueOf(System.currentTimeMillis())); redisTemplate.opsForHash().putAll(taskKey, meta); redisTemplate.expire(taskKey, TASK_TTL); request.setTaskId(taskId); // 2. 异步投递消息到 RocketMQ,保障生产端高吞吐 rocketMQTemplate.asyncSend(TOPIC_LLM_TASK, MessageBuilder.withPayload(request).build(), new SendCallback() { @Override public void onSuccess(SendResult sendResult) { log.info("任务消息投递成功, taskId: {}, msgId: {}", taskId, sendResult.getMsgId()); } @Override public void onException(Throwable throwable) { log.error("任务消息投递失败, taskId: {}", taskId, throwable); redisTemplate.opsForHash().put(taskKey, "status", TaskStatus.SUBMIT_FAILED.name()); } }); // 3. 立即向客户端返回 202 Accepted TaskResponse response = TaskResponse.builder() .taskId(taskId) .status(TaskStatus.PENDING) .message("任务已受理并排队中") .estimatedWaitTimeSec(15) .build(); return ResponseEntity.status(HttpStatus.ACCEPTED).body(response); } @GetMapping("/{taskId}/status") public ResponseEntity<TaskResponse> getStatus(@PathVariable String taskId) { String taskKey = "llm:task:" + taskId; Map<Object, Object> entries = redisTemplate.opsForHash().entries(taskKey); if (entries.isEmpty()) { return ResponseEntity.status(HttpStatus.NOT_FOUND).build(); } String statusStr = (String) entries.get("status"); String result = (String) entries.get("result"); String errorMsg = (String) entries.get("errorMsg"); TaskResponse response = TaskResponse.builder() .taskId(taskId) .status(TaskStatus.valueOf(statusStr)) .result(result) .message(errorMsg) .build(); return ResponseEntity.ok(response); } }Worker 消费端:带流控治理的平滑消费者
消费端的最大挑战在于外部大模型供应商的调用配额限制。如果不加节制地并发拉取,Worker 集群很容易打爆供应商设定的 RPM 阈值。因此,Worker 必须具备平滑的流量整形能力。在单节点维度结合 GuavaRateLimiter,在分布式集群维度结合 Redis 令牌桶或 Sentinel,将外呼并发与频率压制在安全红线以下。
package com.example.worker.consumer; import com.example.gateway.domain.AiTaskRequest; import com.example.gateway.domain.TaskStatus; import com.google.common.util.concurrent.RateLimiter; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.core.RocketMQListener; import org.springframework.ai.chat.client.ChatClient; import org.springframework.data.redis.core.StringRedisTemplate; import org.springframework.stereotype.Component; import java.time.Duration; @Slf4j @Component @RequiredArgsConstructor @RocketMQMessageListener( topic = "LLM_TASK_DISPATCH_TOPIC", consumerGroup = "llm_worker_consumer_group", consumeThreadMax = 8 ) public class AiTaskConsumer implements RocketMQListener<AiTaskRequest> { private final StringRedisTemplate redisTemplate; private final ChatClient.Builder chatClientBuilder; // 速率控制器:单节点限制每秒最多发起 4 次大模型调用,避免瞬时并发触发 429 private final RateLimiter rateLimiter = RateLimiter.create(4.0); @Override public void onMessage(AiTaskRequest request) { String taskId = request.getTaskId(); String taskKey = "llm:task:" + taskId; String lockKey = "llm:lock:" + taskId; // 1. 分布式防重执行:通过 Redis SETNX 获取执行锁,避免网络抖动重试导致多次扣费与重复推理 Boolean acquiredLock = redisTemplate.opsForValue().setIfAbsent(lockKey, "LOCKED", Duration.ofMinutes(10)); if (Boolean.FALSE.equals(acquiredLock)) { log.warn("检测到重复投递或正在执行的任务, taskId: {}", taskId); return; } try { // 2. 消费限流等待 double waitTime = rateLimiter.acquire(); log.debug("获取消费令牌成功, taskId: {}, 等待时长: {}s", taskId, waitTime); // 3. 更新状态机为 PROCESSING redisTemplate.opsForHash().put(taskKey, "status", TaskStatus.PROCESSING.name()); redisTemplate.opsForHash().put(taskKey, "startedAt", String.valueOf(System.currentTimeMillis())); // 4. 调用大模型进行计算 ChatClient chatClient = chatClientBuilder.build(); String aiResult = chatClient.prompt() .user(request.getPrompt()) .call() .content(); // 5. 保存结果并更新状态为 SUCCESS redisTemplate.opsForHash().put(taskKey, "result", aiResult); redisTemplate.opsForHash().put(taskKey, "status", TaskStatus.SUCCESS.name()); redisTemplate.opsForHash().put(taskKey, "completedAt", String.valueOf(System.currentTimeMillis())); log.info("大模型异步任务处理完成, taskId: {}", taskId); } catch (Exception e) { log.error("大模型任务推理失败, taskId: {}", taskId, e); redisTemplate.opsForHash().put(taskKey, "status", TaskStatus.FAILED.name()); redisTemplate.opsForHash().put(taskKey, "errorMsg", e.getMessage()); // 依据业务异常类型决定是否抛出异常以触发 MQ 梯度重试 } finally { redisTemplate.delete(lockKey); } } }生产排坑与高可用治理
在实际生产运营中,仅仅实现基本的消息收发远远不够,必须针对以下复杂异常场景建立兜底与防御机制:
1. 毒丸消息(Poison Pill)防死循环与死信隔离
某些用户的输入可能包含触发大模型安全风控的内容,或者超长 Prompt 导致模型上下文溢出(Context Length Exceeded)。如果直接抛出异常让 MQ 重试,这条消息会在队列中反复拉取、反复报错,消耗宝贵的调用配额并阻塞消费线程。
合理的治理策略是:细分异常类型。对于ModelSecurityException或PromptTooLongException等不可恢复异常,直接将任务状态标记为TERMINATED_BY_POLICY并确认消费;对于SocketTimeoutException或临时429 Too Many Requests,允许按指数退避策略重试,达到最大重试次数(例如 3 次)后自动转入死信队列(DLQ),并触发钉钉或企业微信告警。
2. 多租户与 VIP 优先级队列划分
如果所有业务共用同一个 Topic,当某个批量离线业务突然塞入 10 万条知识库向量化切片任务时,线上核心客户的单条合同审查任务将面临极长的排队延迟。
在消息队列层面,应按租户等级或业务时效性划分独立 Topic,例如LLM_TASK_VIP_TOPIC与LLM_TASK_BATCH_TOPIC。Worker 集群采用差异化线程配比:70% 的消费算力监听 VIP 队列,30% 的算力监听批量队列。当 VIP 队列为空时,Worker 可动态借调算力消费批量队列,保证核心业务的低延迟 SLA。
3. 客户端长轮询与 Webhook 回调联动
为了减少客户端高频轮询给网关和 Redis 带来的 QPS 压力,网关可提供基于 DeferredResult 的长轮询接口(Long-Polling):客户端发起状态查询时,若任务仍处于 PROCESSING,网关挂起请求 15 秒;一旦任务完成,通过 Redis Pub/Sub 唤醒并立即响应。对于耗时超过 5 分钟的超长任务,强烈建议在提交时传入callbackUrl,Worker 处理完成后发起带有重试机制的 HTTP POST 回调,彻底消除无谓的轮询流量。
4. 关键监控与水位预警指标
异步任务池的稳定性高度依赖监控系统的可观测性。在 Prometheus 中必须固化以下核心指标:
- MQ Lag 水位:按 Topic 和 Consumer Group 监控消息堆积量,堆积阈值超过 1000 时触发扩容告警。
- 任务端到端 P90/P99 耗时:从客户端提交到最终结果写入的整体生命周期耗时。
- 上游配额消耗率:实时统计每分钟向大模型厂商发起的请求数(RPM)和 Token 消耗量(TPM),与厂商购买的配额上限做实时比例比对,在达到 85% 水位时自动触发 Worker 端的消费降速,实现闭环的主动防御。