ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Spring Boot 3 异步生图管道与防刷架构实战

Spring Boot 3 异步生图管道与防刷架构实战 1. 为什么要在 Spring Boot 3 里做 AI 生图管道1.1 从一次线上事故说起去年年底我接手了一个内部创意工具的后端核心功能是让运营同学输入一段中文描述调用生图模型返回一张配图。最初版本极其简陋Controller 里直接RestTemplate调模型接口同步等待拿到 URL 就返回。上线第一周就出事了——晚高峰时段十几个运营同时点“生成”Tomcat 线程池瞬间被打满整个服务连带其他接口一起 502。更糟的是有个同事写了个脚本循环调用一晚上烧掉了小两千的额度。那次事故之后我花了大概三周时间把这套东西彻底重构成了一条异步生图管道并且加上了防刷架构。这篇文章就是把这套方案的完整思路和落地细节摊开讲一遍。如果你正在用 Spring Boot 3 对接生图类模型或者准备做类似的 AI 能力中台这篇内容应该能帮你少踩不少坑。先说清楚这套东西是什么、能干什么、适合谁看。它本质上是一个后端服务层的设计模式把“接收请求 → 参数校验 → 排队 → 调用模型 → 存储结果 → 回调通知”这一整条链路拆开用异步管道串起来再叠加限流、幂等、额度控制等防刷手段。适合有一定 Spring Boot 基础、正在做 AI 能力集成的后端同学也适合想了解工业级 AI 接口怎么设计的架构方向读者。哪怕你用的是别的语言栈管道和防刷的思路是通用的。1.2 同步调用为什么必然翻车很多人第一反应是“生图不就是调个接口吗为什么要搞这么复杂”。问题就出在生图这个动作的时间特性上。文本类接口通常几百毫秒返回而生图接口普遍在 5 到 30 秒之间遇到排队高峰甚至更久。这个量级的时间差决定了同步调用在生产环境里几乎必然出问题。我列一下同步方案的核心痛点你可以对照自己的项目看看中了几条线程占用时间长一个请求占一个 Tomcat 线程几十秒并发稍微上来线程池就满了这是最致命的。超时边界难定网关超时、Nginx 超时、客户端超时三层超时如果没对齐会出现“模型其实生成成功了但用户看到失败”的尴尬情况钱花了图没了。无法重试同步链路里一旦失败重试逻辑很难写因为用户已经等了很久了。防刷无从下手请求和计算强绑定你没法在“排队阶段”做削峰只能硬扛。结果无法复用同样的 prompt 每次都要重新生成浪费额度。异步管道恰好能把这几个问题逐个化解。请求进来先落库、入队立刻返回一个任务 ID用户拿着 ID 去轮询或等推送。计算资源和请求流量解耦之后削峰、重试、幂等、缓存全都变得可做了。1.3 整体架构的分层设计我把整套系统分成四层从上到下依次是接入层、调度层、执行层、存储层。这个分层不是拍脑袋定的每一层都对应一个明确的职责边界边界清晰了后面加功能才不会互相污染。接入层负责 HTTP 入口、鉴权、参数校验、限流。这一层要尽可能轻只做“接住请求”这件事绝不碰业务逻辑。调度层负责把请求转成任务、写入队列、管理任务状态机。执行层是真正调用生图模型的地方也是唯一会长时间阻塞的地方它独立部署、独立扩容。存储层负责图片落盘、元数据入库、结果缓存。这么分的好处是执行层挂了不影响接入层收请求接入层扩容不影响执行层的并发控制。你可以把执行层想象成工厂车间接入层是前台接待前台只管收单子车间按自己的节奏生产中间靠队列这个“传送带”连接。提示分层的关键是依赖单向。接入层可以依赖调度层调度层可以依赖执行层但绝不能反向依赖。我见过有人图省事让执行层直接回调接入层的 Controller结果循环依赖加事务嵌套排查了两天。2. 核心组件的选型与取舍2.1 队列选型为什么最终选了 Redis Stream队列是整条管道的血管选型必须慎重。我评估过三个方案数据库轮询表、RabbitMQ、Redis Stream。最后选了 Redis Stream理由如下。数据库轮询表是最容易想到的方案一张task表状态字段标记待处理执行层定时select ... for update捞任务。优点是简单、事务一致性好、不用引入新组件。缺点是轮询有延迟、高频轮询压数据库、并发抢锁容易出问题。适合任务量很小的场景日任务量上千以内可以凑合。RabbitMQ 是成熟方案功能全、有死信队列、有 ACK 机制。缺点是运维成本高多一个中间件就多一份负担而且它的消息确认模型和我们的“任务状态机”有点重叠容易搞出两套状态。Redis Stream 是我最终的选择。它有几个特别契合的点消费者组天然支持多实例竞争消费执行层可以水平扩容消息持久化 ACK 机制保证任务不丢XPENDING和XCLAIM能处理消费者宕机后的消息转移这是做可靠队列的关键而且大部分项目本来就有 Redis不用额外引入组件。# 创建消费者组从最新消息开始消费 XGROUP CREATE image:stream image:group $ MKSTREAM # 消费者读取消息一次读10条阻塞2秒 XREADGROUP GROUP image:group consumer-1 COUNT 10 BLOCK 2000 STREAMS image:stream 注意XADD时一定要用MAXLEN ~ 10000做近似裁剪否则 Stream 会无限增长把内存吃满。~号是近似裁剪性能比精确裁剪好很多代价是长度可能略微超标这个误差完全可以接受。2.2 状态机设计任务生命周期怎么管任务状态机是防刷和幂等的基石。我定义了六个状态PENDING已入队、PROCESSING生成中、SUCCESS成功、FAILED失败、REJECTED被限流拒绝、EXPIRED超时作废。状态流转必须严格受控我用一张表把允许的流转列清楚代码里用枚举校验非法流转直接抛异常。这样做的价值在于任何时刻你都能通过状态字段准确知道任务卡在哪一步排查问题不用猜。当前状态允许流转到触发条件PENDINGPROCESSING / REJECTED / EXPIRED被消费 / 限流 / 超时未消费PROCESSINGSUCCESS / FAILED模型返回 / 模型报错或重试耗尽FAILEDPENDING人工或自动重试SUCCESS终态不可再变REJECTED终态不可再变这里有个容易忽略的点状态更新必须带乐观锁。执行层可能因为网络抖动重复消费同一条消息如果不用版本号或where status PENDING这样的条件更新就会出现同一任务被处理两次、扣两次额度的情况。我的做法是每次更新都带上原状态作为条件update task set statusPROCESSING where id? and statusPENDING影响行数为 0 就说明被别人抢先了直接跳过。2.3 幂等设计同一请求只扣一次额度幂等的核心是给每个请求一个唯一标识。我让客户端在请求头里带一个X-Request-Id服务端用它做幂等键。如果客户端没带就用“用户ID prompt 的哈希 时间窗口”兜底生成一个。具体实现是入队前先查 Redis 里有没有这个幂等键有就直接返回上次的任务 ID不重复入队、不重复扣额度。这个键的过期时间设成 10 分钟足够覆盖用户手抖连点的情况又不会长期占用内存。String idempotentKey idem: userId : requestId; Boolean first redisTemplate.opsForValue() .setIfAbsent(idempotentKey, taskId, Duration.ofMinutes(10)); if (Boolean.FALSE.equals(first)) { // 重复请求直接返回已有任务ID return existingTaskId; }实操心得幂等键一定要在扣额度之前判断。我第一版写反了顺序先扣额度再判幂等结果用户连点两次扣了两份额度但只生成一张图被投诉了好几次。顺序错了逻辑再对也没用。3. 异步管道的完整落地实现3.1 接入层请求进来先过三道关接入层的 Controller 我写得非常克制整个方法体不超过 20 行。它只做三件事鉴权、参数校验、限流判断然后交给调度层。第一道关是鉴权。用 Spring Security 的过滤器链拿到用户身份塞进SecurityContext。这一步是后面所有防刷逻辑的基础没有用户身份就没法做按用户的额度控制。第二道关是参数校验。用 Jakarta Validation 注解校验 prompt 长度、图片尺寸、数量等。这里有个细节prompt 长度上限不要设太大我设的是 800 字符。太长的 prompt 不仅浪费 token还容易被用来做注入攻击。校验不通过直接返回 400不进入后续流程。第三道关是限流。我用的是基于 Redis 的滑动窗口限流按用户维度限制“每分钟 N 次、每天 M 次”。限流不通过返回 429并且带上Retry-After头告诉客户端多久后重试。PostMapping(/images/generate) public ResultTaskVO generate(Valid RequestBody GenerateRequest req) { Long userId SecurityUtils.getCurrentUserId(); // 滑动窗口限流每分钟10次 if (!rateLimiter.tryAcquire(userId, 10, Duration.ofMinutes(1))) { throw new BizException(ErrorCode.RATE_LIMITED); } // 每日额度检查 quotaService.checkAndDeduct(userId, req.getCount()); return Result.ok(scheduler.submit(userId, req)); }3.2 调度层任务落库与入队的一致性调度层最容易出问题的地方是数据库写入和队列入队的一致性。如果先写库再入队入队失败任务就永远卡在 PENDING如果先入队再写库执行层可能拿到一个数据库里还不存在的任务 ID。我的解法是先写库、再入队、入队失败标记任务为 FAILED 并告警。因为写库是本地事务成功率高入队失败是极小概率事件标记失败后由补偿任务兜底重试。这个顺序保证了不会出现“队列里有、库里没有”的幽灵任务。Transactional public TaskVO submit(Long userId, GenerateRequest req) { Task task new Task(); task.setUserId(userId); task.setPrompt(req.getPrompt()); task.setStatus(TaskStatus.PENDING); taskMapper.insert(task); // 事务提交后再入队避免事务未提交就被消费 TransactionSynchronizationManager.registerSynchronization( new TransactionSynchronization() { Override public void afterCommit() { try { streamProducer.send(task.getId()); } catch (Exception e) { taskMapper.markFailed(task.getId(), enqueue failed); alertService.warn(入队失败: task.getId()); } } }); return TaskVO.from(task); }注意入队一定要放在afterCommit里。我踩过的坑是直接在事务方法里入队结果执行层消费太快去数据库查任务时事务还没提交查不到记录任务被当成脏数据丢弃了。这个 bug 在低并发时根本复现不出来高并发才偶发排查了很久。3.3 执行层消费、调用、回写三步走执行层是独立部署的服务启动时注册消费者组循环从 Stream 里读消息。每读一条走“消费 → 调用模型 → 回写结果”三步。消费阶段先做状态抢占用条件更新把任务从 PENDING 改成 PROCESSING抢不到就说明被别人处理了直接 ACK 掉这条消息。抢占成功后开始调用模型。调用模型这块要重点说超时和重试。生图接口的超时我设的是 60 秒比模型平均耗时长一倍留足余量。重试策略是“最多重试 2 次指数退避”第一次失败等 2 秒第二次等 4 秒。重试只针对网络超时和 5xx 错误4xx 参数错误不重试因为重试也没用。StreamListener public void consume(StreamMessage message) { Long taskId Long.valueOf(message.getBody()); // 状态抢占 int updated taskMapper.casStatus(taskId, PENDING, PROCESSING); if (updated 0) { streamConsumer.ack(message); return; } try { String imageUrl imageClient.generate(task.getPrompt(), 60); taskMapper.markSuccess(taskId, imageUrl); } catch (RetryableException e) { retryService.schedule(taskId, e.getAttempt()); } catch (Exception e) { taskMapper.markFailed(taskId, e.getMessage()); } finally { streamConsumer.ack(message); } }回写结果时图片要先上传到对象存储拿到永久 URL 再写库。千万别直接把模型返回的临时 URL 存库那些 URL 通常几小时就失效了用户过两天再看图就 404 了。这个坑我替你们踩过了。3.4 存储层图片落盘与元数据分离存储层我做了图片和元数据分离。图片本体放对象存储数据库只存 URL、prompt、参数、用户 ID、生成耗时这些元数据。这样做的好处是数据库轻量查询快图片的存储和 CDN 加速交给专业组件。元数据表我加了几个关键索引user_id created_at用于查用户历史status created_at用于监控和补偿任务扫描idempotent_key唯一索引用于幂等。索引不是越多越好每加一个都要想清楚它服务于哪个查询否则写入性能会被拖垮。另外我加了一张额度流水表每次扣额度都记一条包含用户 ID、变动数量、变动原因、关联任务 ID。这张表是排查“为什么用户说额度不对”的终极武器有它在任何额度争议都能三分钟查清楚。4. 防刷架构的层层设防4.1 四层限流从粗到细的漏斗防刷不能只靠一层限流我用的是四层漏斗从粗到细逐层过滤。第一层是网关层限流按 IP 限制总 QPS挡住最粗暴的脚本攻击。这一层阈值设得比较宽只挡明显异常的流量。第二层是用户维度限流就是我前面说的滑动窗口按用户限制每分钟和每天的调用次数。这是最核心的一层因为大部分滥用都来自正常账号。第三层是全局并发控制限制执行层同时处理的任务数。用 Redis 的计数器实现超过阈值的新任务直接排队等待而不是拒绝。这一层保护的是模型接口防止我们把上游打挂。第四层是成本熔断按天统计总消耗超过预算阈值就自动降级只允许白名单用户调用。这一层是最后的保险丝防止出现“一晚上烧两千”的事故重演。层级维度阈值示例超限动作网关层IP100 QPS拒绝用户层用户ID10次/分、200次/天拒绝并发层全局50 并发排队成本层全局日预算上限降级4.2 行为风控识别机器与羊毛党光靠频率限流还不够有些羊毛党会用多个账号、控制频率来绕过。这时候需要行为风控。我采集了几个关键特征请求间隔的方差、prompt 的相似度、调用时间的分布。正常用户的请求间隔是随机的方差大脚本调用的间隔非常规律方差极小。正常用户的 prompt 各不相同脚本往往用同一批 prompt 反复刷。正常用户白天活跃脚本可能 24 小时不停。基于这些特征我给每个用户算一个风险分超过阈值就触发二次验证或直接限流。这套东西不需要多复杂的模型几个简单的统计规则就能挡住 90% 的脚本。实操心得风控规则一定要留观察期。新规则上线先只记录不拦截观察一周看误伤率确认没问题再开启拦截。我第一版风控规则太激进把一个正常的高频用户给封了人家是设计团队在做批量出图差点闹到领导那里。4.3 额度体系预扣、返还与对账额度体系的设计要点是预扣 返还。请求进来先预扣额度生成成功确认扣除生成失败自动返还。这样既防止了并发超额又保证了失败不扣钱。预扣用 Redis 的原子操作实现DECRBY之后如果结果小于 0 就回滚。这里要注意预扣和返还必须成对出现我在代码里用 try-finally 保证返还逻辑一定执行哪怕中间抛异常。对账是每天凌晨跑的定时任务把额度流水表的变动汇总和用户当前余额比对不一致就告警。这个对账机制帮我抓到过一次并发 bug两个请求同时预扣因为没用原子操作导致额度扣少了。没有对账这种 bug 可能几个月都发现不了。4.4 内容安全prompt 与生成结果的双重过滤内容安全这块不能省。我在两个环节做了过滤prompt 入队前和图片生成后。prompt 过滤用的是敏感词库 正则规则命中就直接拒绝不消耗额度。词库要定期更新我把它放在配置中心改完实时生效不用重启服务。图片过滤是在生成成功后、返回给用户前调用内容审核接口过一遍。审核不通过的图片直接删除任务标记为 FAILED并且记录用户行为。这里有个权衡审核会增加延迟所以我把它放在异步链路里用户拿到的是“审核中”状态审核通过后才变成“成功”。注意内容审核接口也可能超时或失败。我的策略是审核失败时默认拦截宁可误伤不可放过。同时记录日志人工复核。安全问题上保守永远比激进好。5. 常见问题与排查实录5.1 任务卡在 PROCESSING 不动了这是最常见的问题。原因通常有三个执行层实例挂了、模型接口长时间无响应、回写数据库失败。排查顺序是先看执行层实例的健康状态和日志再看模型接口的监控指标最后查数据库里这些任务的状态更新时间和重试次数。如果是实例挂了Stream 的XPENDING会显示消息未 ACK重启实例后会自动重新消费。如果是模型无响应超时机制会触发重试。如果是回写失败通常是数据库连接池耗尽看连接池监控就能定位。我加了一个僵尸任务扫描定时任务每 5 分钟扫一次把 PROCESSING 状态超过 10 分钟的任务重新入队。这个兜底机制救过我好几次。5.2 用户反馈“扣了额度没出图”这类问题九成是幂等或状态流转的 bug。排查步骤先用任务 ID 查任务状态再用用户 ID 查额度流水把两条时间线对齐看。常见的情况是任务实际成功了但客户端因为网络问题没收到结果用户以为失败了又点了一次第二次被幂等拦截返回了第一次的任务 ID但客户端没处理好这个响应。这种问题要在客户端配合下解决服务端能做的就是保证幂等键和任务 ID 的映射稳定。还有一种情况是预扣了但返还逻辑没执行通常是代码里 try-finally 写漏了或者返还时抛了异常被吞掉。所以返还逻辑一定要加日志和告警。5.3 高峰期任务积压严重积压说明消费速度跟不上生产速度。先看是生产太快还是消费太慢。生产太快就加强限流消费太慢就扩容执行层实例。扩容执行层时要注意消费者组的并发上限。Redis Stream 的消费者组里同一个消费者名只能有一个实例多个实例要用不同的消费者名。我一开始所有实例用同一个消费者名结果只有一个实例在消费其他都在空转白白浪费资源。另外模型接口本身可能有并发限制扩容到一定程度就上不去了。这时候要考虑分级队列付费用户走高速队列免费用户走普通队列保证核心用户体验。5.4 问题速查表现象可能原因排查动作解决方式任务卡 PROCESSING实例挂/模型超时/回写失败查实例状态、模型监控、DB连接池重启实例/等重试/扩容连接池扣额度没出图幂等bug/返还漏执行对齐任务与流水时间线修幂等逻辑/补返还任务积压消费慢/生产快看生产消费速率对比扩容/加强限流/分级队列重复扣额度幂等键失效/并发预扣查幂等键TTL、预扣原子性修TTL/改原子操作图片404存了临时URL查存储层URL来源改为上传对象存储后存永久URL6. 一些踩坑之后的经验之谈6.1 关于超时对齐这件事超时对齐是我认为最容易被忽视、但影响最大的细节。客户端超时、网关超时、服务端处理超时、模型调用超时这四个超时必须形成递增的阶梯客户端 网关 服务端 模型调用。如果客户端超时比服务端还短用户会在服务端还在处理时就收到超时然后重试造成重复任务。我的配置是客户端 90 秒、网关 80 秒、服务端异步接口 5 秒返回任务 ID、模型调用 60 秒。注意异步接口本身是秒回的所以服务端的 5 秒超时是针对“入队”这个动作不是针对整个生成过程。这个区分很关键很多人把异步接口的超时也设成 60 秒那就失去异步的意义了。6.2 关于日志与可观测性异步链路的排查难度远高于同步链路因为请求和响应不在一个线程里。所以全链路追踪是必须的。我在任务 ID 生成时就把它塞进 MDC后续所有日志都带上这个 ID包括执行层的日志。这样用任务 ID 一搜整条链路的日志全出来了。除了日志还要有指标监控。我埋了几个关键指标入队速率、消费速率、成功率、平均耗时、各状态任务数。这些指标画成看板一眼就能看出系统是否健康。有一次成功率突然从 98% 掉到 85%看板立刻报警查下来是模型接口那边在灰度及时切了备用通道。6.3 关于成本控制生图是实打实花钱的成本控制必须从第一天就做。我的做法是每个任务都记录消耗按用户、按天、按 prompt 类型多维统计。这样你能清楚知道钱花在哪了。省钱的一个实用技巧是结果缓存。同样的 prompt 同样的参数直接返回缓存结果不重新生成。缓存键用 prompt 和参数的哈希缓存有效期设 7 天。我们内部工具里重复 prompt 的比例高达 30%光这一项就省了不少。另一个技巧是参数降级。非核心场景用低分辨率、少步数的参数成本能降一半以上。用户如果对质量有要求再手动升级参数。这个策略要在产品层面和用户说清楚避免预期落差。6.4 关于灰度与回滚任何改动都要能灰度、能回滚。我的做法是用配置开关控制新逻辑比如新的限流规则、新的重试策略都放在配置中心改完实时生效。出问题一键关掉回到旧逻辑。执行层的版本升级用滚动发布一次只更新一个实例观察几分钟没问题再更新下一个。因为执行层是无状态的状态都在数据库和 Redis滚动发布很安全。但要注意发布期间正在处理的任务会中断靠僵尸任务扫描兜底重新入队。这套系统跑了大半年从最初的每天几十个任务到现在日均几千个中间经历过几次流量高峰都没出大问题。回头看最有价值的不是某个具体的技术选型而是分层解耦 状态机 幂等 多层限流这套组合拳的思路。技术会变模型会换但这套设计模式是通用的。你要是也在做类似的东西建议先把状态机和幂等这两块打扎实剩下的都是水到渠成的事。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进