ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

大模型网关削峰填谷:基于消息队列的异步任务池

大模型网关削峰填谷:基于消息队列的异步任务池 大模型网关削峰填谷基于消息队列的异步任务池在大模型系统落地到实际生产业务的过程中除了智能对话、搜索补全这类必须在数百毫秒内返回的首字交互场景外有大量业务天然属于长耗时、高计算密度的离线或半离线任务。典型的场景包括企业级合同长文档合规审查、全量代码库安全扫描与重构建议、数十万字营销文案的批量生成以及企业私域知识库的大规模切片与向量化嵌入。这类任务在单次推理或多次链式调用Agent Loop中耗时往往从数十秒延伸至数分钟。如果网关层依然沿用传统的同步 HTTP 或 RPC 调用模式客户端保持长连接阻塞等待系统的网络文件描述符FD、Tomcat 工作线程池以及微服务连接池会在短时间内被耗尽。更致命的是一旦上游业务在某一时刻批量提交成百上千个文档突发的并发流量会瞬间击穿网关向上游大模型服务商发起雪崩式请求触发大面积的 HTTP 429 Too Many Requests 错误导致整体业务瘫痪。解决这一工程痛点的标准解法是在大模型网关中引入基于消息队列Message Queue的异步任务池架构通过“异步接收、状态落库、消息缓冲、受控消费、结果通知”的全流程闭环实现流量的高效削峰填谷。架构设计从同步阻塞到异步流水线在异步削峰架构中核心思想是将客户端的“任务提交请求”与底层的“推理计算执行”彻底解耦。整个架构由四个核心组件构成接入网关API Gateway负责鉴权、参数校验、生成全局唯一任务 ID、初始化任务状态机并将任务载荷投递到消息队列随后立即向客户端返回 HTTP 202 Accepted 状态码及任务凭证。消息中枢RocketMQ / Kafka作为削峰填谷的蓄水池承载瞬时峰值流量隔离前后端处理速度的不对称性提供可靠的消息持久化和顺序保证。AI Worker 计算集群作为消息消费者根据预设的 RPMRequests Per Minute和 TPMTokens 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 ResponseEntityTaskResponse submitTask(RequestBody AiTaskRequest request) { String taskId TASK_ UUID.randomUUID().toString().replace(-, ); String taskKey llm:task: taskId; // 1. 初始化任务状态机写入 Redis MapString, 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 ResponseEntityTaskResponse getStatus(PathVariable String taskId) { String taskKey llm:task: taskId; MapObject, 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 RocketMQListenerAiTaskRequest { 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 分钟的超长任务强烈建议在提交时传入callbackUrlWorker 处理完成后发起带有重试机制的 HTTP POST 回调彻底消除无谓的轮询流量。4. 关键监控与水位预警指标异步任务池的稳定性高度依赖监控系统的可观测性。在 Prometheus 中必须固化以下核心指标MQ Lag 水位按 Topic 和 Consumer Group 监控消息堆积量堆积阈值超过 1000 时触发扩容告警。任务端到端 P90/P99 耗时从客户端提交到最终结果写入的整体生命周期耗时。上游配额消耗率实时统计每分钟向大模型厂商发起的请求数RPM和 Token 消耗量TPM与厂商购买的配额上限做实时比例比对在达到 85% 水位时自动触发 Worker 端的消费降速实现闭环的主动防御。
返回列表