1. 项目背景与整体设计思路拆解
先说个大家都懂的现状:AI 生图接口的接入,从代码层面看就是一个 HTTP 调用,POST 一段 JSON 进去,拿回一张图片或者一个 URL,完事了。但真正把这件事放到生产环境里,你会发现“能跑”和“能上线”之间隔着的是一条相当宽的河。
我在几个项目里接过 gpt-image-2.5 这类生图模型的 API,也帮团队搭建过完整的生图服务,感受最深的一点是:图片生成本身只占整个工作量的两到三成,剩下七八成的精力全都耗在管道稳定性和防刷治理上。你想想,你自己的账号、自己的 Key,如果直接被前端拿去调用,那意味着你的成本、你的接口、你的配额全部暴露在公网上。任何人拿到你的页面地址,就可以无限刷图,轻则账单爆炸,重则接口被恶意调用打到限流甚至封禁。这不是危言耸听,是真实发生过的生产事故。
所以这篇文章,我想用 Spring Boot 3 做载体,完整拆解一条工业级 AI 生图管道的搭建过程。涉及到的东西包括:OpenAI 兼容接口的客户端封装、任务队列与状态机设计、异步处理与重试策略、分布式限流、用户额度控制、幂等防重,以及最后的压测验证。内容偏实操,代码为主,每一个关键决策我都会讲清楚为什么这么做。
1.1 从 Demo 到工业级的差距:三个必须补上的短板
如果你只是本地写个 Demo 玩,调用 gpt-image-2.5 的流程简单到令人发指:一个 RestTemplate 或者 OkHttp 请求,加上 API Key 头,参数一填,图片就回来了。但一旦面对真实用户流量,你会发现这套直连模式有三个绕不过去的短板。
第一个短板是耗时不可控。生图接口不是普通 REST 接口,它通常需要十几秒甚至几十秒才能返回结果。HTTP 连接池的线程如果被长时间占用,你的服务吞吐量会急剧下降。Tomcat 默认 200 个线程,假设平均每个生图请求占用 20 秒,那这个服务的 QPS 天花板就是 10,而且这 10 个请求还把所有线程都占满了,其他接口全部跟着遭殃。所以必须把生图调用从请求线程中剥离出去,用异步任务的方式处理。
第二个短板是失败率不低。模型服务端的负载、网络抖动、超时、限流,任何一个环节出问题都会导致生图失败。直连模式下用户看到的就是一个 500 或者超时,体验极差。工业级管道需要内部消化这些异常:重试、熔断、降级、失败通知,每一步都要有兜底。
第三个短板也是我今天重点讲的,成本与安全。生图 API 是明码标价按张计费的,这意味着每一个请求都在消耗真金白银。如果不对调用方做身份认证、额度控制、限流和防刷,那你的服务就是一个敞开的 ATM 机。真实项目里我见过有人拿别人的生图接口做“免费图片生成工具”往外卖流量的,也见过内部测试 Key 被截获后被人拿来刷了上万张图的。防刷不是可选项,是必选项。
1.2 整体架构分层:一条可复用的流水线模型
我最终采用的架构可以抽象成五层,每一层职责单一,层与层之间通过接口解耦。这个分层模型不绑定具体模型,以后从 gpt-image-2.5 换成别的生图模型,只需要替换最底层的适配器。
第一层是接入层,负责接收用户的生图请求。包括参数校验、身份认证、幂等判断、限流拦截。这一层的目标是:用最小的成本挡掉绝大多数无效流量和恶意流量。
第二层是任务管理层,核心是任务状态机和调度队列。生图请求到达后立即返回一个任务 ID,用户通过任务 ID 查询进度和结果。异步执行器从队列中消费任务,更新状态、执行重试、处理回调。
第三层是模型适配层,也就是真正调用 gpt-image-2.5 的地方。这一层做统一的请求封装、响应解析、错误映射。对外暴露的是一个简单的generateImage(request)方法,内部屏蔽了 HTTP 细节。
第四层是存储层,包括任务数据、用户额度数据、图片文件/URL 的存储。我用的组合是 MySQL + Redis,MySQL 管任务和用户持久化数据,Redis 管实时计数、分布式锁和令牌桶。
第五层是治理层,监控、日志、告警、配额管理。这一层不直接参与业务逻辑,但决定了整个服务能不能长期稳定运行。
这五层的关系,你可以理解成一条自动化的流水线:接入层是安检门,任务管理层是传送带,模型适配层是加工车间,存储层是仓库,治理层是监控室。每一层都会在下面的代码实现串起来。
2. Spring Boot 3 项目搭建与基础客户端封装
工欲善其事,必先利其器。Spring Boot 3 相比 2.x 有一些底层变化,最需要关注的是两点:一是 Jakarta EE 9+ 的包名迁移,javax.*变成了jakarta.*;二是 Spring 6 引入的RestClient,这是一个比 RestTemplate 更现代化的 HTTP 客户端,接口风格流畅,支持同步和异步,非常适合用来封装对外的 API 调用。我下面的代码全部基于 Spring Boot 3.2.x。
2.1 项目初始化和基础依赖配置
创建工程我用的是 Spring Initializr,选择 Java 17 + Spring Boot 3.2.5。依赖方面,除了 Web 和 Validation,我额外加了这几个:Spring Data Redis(用于限流和分布式锁)、MySQL 驱动和 MyBatis-Plus(用于任务和用户数据的持久化),以及 Spring Boot Actuator(用于健康检查和监控指标)。Lombok 省去样板代码。
spring.application.name=image-pipeline-service server.port=8080 # 连接池 spring.datasource.url=jdbc:mysql://localhost:3306/image_pipe spring.datasource.username=root spring.datasource.password=yourpassword spring.datasource.hikari.maximum-pool-size=20 # Redis spring.data.redis.host=localhost spring.data.redis.port=6379 spring.data.redis.timeout=3s # 生图模型配置 gpt-image.api-key=${IMAGE_API_KEY} gpt-image.base-url=https://api.example.com/v1 gpt-image.model=gpt-image-2.5 gpt-image.timeout-seconds=60需要注意的一点是,API Key 一定不要写在配置文件里提交到 Git,用环境变量或者配置中心注入是底线。我见过不止一次有人把 Key 硬编码在 application.yml 里,然后整个仓库被拉下来抄走,直接导致账号被盗刷。这不是技术问题,是安全意识问题。
线程池配置这里,我用的是 Spring Boot 3.2 支持的虚拟线程。对于生图这种 IO 密集型的任务,虚拟线程的性价比远高于平台线程。但也要说明,虚拟线程在 JDBC 等同步阻塞场景下使用要谨慎,好在我们操作 MySQL 是短暂的,风险可控。
spring.threads.virtual.enabled=true如果你对虚拟线程还不够放心,也可以用传统线程池,后续代码里我会给出两种写法,你根据自己的场景选。
2.2 封装 gpt-image-2.5 的 OpenAI 兼容客户端
gpt-image-2.5 的接口格式遵循 OpenAI Images API 的风格,最简单的文生图请求是这样的:
POST /v1/images/generations Authorization: Bearer sk-xxx Content-Type: application/json { "model": "gpt-image-2.5", "prompt": "A cute corgi sitting on the grass, sunset lighting", "n": 1, "size": "1024x1024" }响应格式也基本固定,返回的是图片的 URL 或者 base64 数据:
{ "created": 1735000000, "data": [ { "url": "https://example.com/images/xxxx.png" } ] }基于这个格式,我做了三层封装。第一层是请求和响应 DTO,第二层是一个ImageClient接口,第三层是使用RestClient的具体实现。
DTO 层面很简单,请求体我用一个 record 定义,方便也安全:
public record ImageGenerationRequest( String model, String prompt, Integer n, String size ) {}响应体我用两个 record 嵌套,一个是单个图片结果,一个是整体响应:
public record ImageData(String url, String b64_json) {} public record ImageResponse(Long created, List<ImageData> data) {}核心的客户端实现,我用了 Spring 6 的RestClient。这里有一个细节要注意:RestClient 默认的超时时间是无限的,你必须显式配置连接超时和读取超时,否则服务端挂起时你的线程会被一直占用。
@Configuration public class ImageClientConfig { @Value("${gpt-image.base-url}") private String baseUrl; @Value("${gpt-image.timeout-seconds}") private int timeoutSeconds; @Bean public RestClient imageRestClient() { var factory = ClientHttpRequestFactorySettings.defaults() .withConnectTimeout(Duration.ofSeconds(5)) .withReadTimeout(Duration.ofSeconds(timeoutSeconds)); return RestClient.builder() .baseUrl(baseUrl) .requestFactory(new JdkClientHttpRequestFactory(factory)) .build(); } }这里我用的是 JDK 自带的 HttpClient 工厂,你也可以换成 Apache HttpClient 或者 OkHttp,效果差别不大。关键是连接超时 5 秒、读取超时 60 秒这个组合,你要根据实际模型响应时间做调整,标准是:读取超时略大于模型 P99 响应时间,太短容易误杀慢请求,太长会让故障恢复变得很慢。
然后是真正的调用实现。我把原来的同步调用和重试逻辑分开写,这里先看基本版:
@Service public class GptImageClient { private final RestClient restClient; @Value("${gpt-image.api-key}") private String apiKey; @Value("${gpt-image.model}") private String model; public GptImageClient(@Qualifier("imageRestClient") RestClient restClient) { this.restClient = restClient; } public ImageData generate(String prompt, String size, int n) { ImageGenerationRequest request = new ImageGenerationRequest(model, prompt, n, size); ImageResponse response = restClient.post() .uri("/images/generations") .header("Authorization", "Bearer " + apiKey) .contentType(MediaType.APPLICATION_JSON) .body(request) .retrieve() .body(ImageResponse.class); if (response == null || response.data() == null || response.data().isEmpty()) { throw new ImageGenerationException("empty response from upstream"); } return response.data().get(0); } }这个类看起来简单,但有几个细节我特意处理了。Bearer前缀是 OpenAI 兼容接口的固定要求,漏了会直接 401。响应里 data 数组可能为空,要主动判空而不是直接取第一个元素。另外,如果你的接口返回的是b64_json而不是 URL,那你需要在返回前把 base64 解码成图片文件存到自己的 OSS 或者本地磁盘,这个逻辑我放到存储环节处理。
2.3 错误映射与异常体系设计
调用第三方接口,错误处理的好坏直接决定排查问题的效率。我把模型接口的 HTTP 错误码映射成了领域异常类,这样业务代码里就不需要关心 HTTP 细节。
public class ImageGenerationException extends RuntimeException { private final int upstreamStatusCode; // constructor, getter... }错误映射的逻辑写在客户端实现里,通过onStatus回调处理:
public ImageData generateWithMapping(String prompt, String size, int n) { ImageGenerationRequest request = new ImageGenerationRequest(model, prompt, n, size); return restClient.post() .uri("/images/generations") .header("Authorization", "Bearer " + apiKey) .contentType(MediaType.APPLICATION_JSON) .body(request) .retrieve() .onStatus(HttpStatusCode::is4xxClientError, (req, res) -> { throw new ImageUpstreamException("upstream 4xx: " + res.getStatusCode(), res.getStatusCode().value()); }) .onStatus(HttpStatusCode::is5xxServerError, (req, res) -> { throw new ImageUpstreamException("upstream 5xx: " + res.getStatusCode(), res.getStatusCode().value()); }) .body(ImageResponse.class) .data() .get(0); }这样设计之后,上游返回 401(Key 失效)、429(触发限流)、500(服务端异常),到了业务层就是不同类型的异常,可以分别配置重试策略和报警规则。尤其是 429,意味着你的频率可能已经接近上游配额上限,这往往是设计预警系统的关键信号。
3. 工业级生图管道的核心实现
客户端封装好了,但这只是管道的第一段。接下来要解决的是任务模型、异步执行、重试策略和结果存储这四个核心环节。我一个个拆开讲。
3.1 任务状态机设计:一张图看清任务的完整生命周期
生图任务从提交到完成,中间会经历多个状态。如果只是在内存里用一个布尔值标记“完成/未完成”,后续排查问题和做补偿都会非常痛苦。我惯用的状态枚举有七个:CREATED、QUEUED、PROCESSING、SUCCEEDED、FAILED、CANCELLED、TIMEOUT。
状态流转我贴个简单的逻辑说明:
- 请求进来先创建任务,状态为
CREATED,写入数据库。 - 任务入队后变更为
QUEUED。 - 消费者拿到任务开始调用模型,状态为
PROCESSING。 - 生成成功,状态为
SUCCEEDED,并把图片 URL 写入结果字段。 - 重试次数耗尽仍然失败,状态为
FAILED。 - 用户主动取消,状态为
CANCELLED。 - 任务在队列中停留超过最大等待时间,状态为
TIMEOUT。
数据库表结构里我用一个status字段加一个version字段做乐观锁,防止并发更新。表的简化 DDL 如下:
CREATE TABLE image_task ( id BIGINT PRIMARY KEY AUTO_INCREMENT, task_no VARCHAR(64) NOT NULL UNIQUE, user_id BIGINT NOT NULL, prompt TEXT NOT NULL, size VARCHAR(32) NOT NULL, status VARCHAR(20) NOT NULL, image_url VARCHAR(512), failure_reason VARCHAR(512), retry_count INT DEFAULT 0, version INT DEFAULT 0, created_at DATETIME NOT NULL, updated_at DATETIME NOT NULL, KEY idx_user_status (user_id, status), KEY idx_status_created (status, created_at) );task_no是业务侧的任务编号,用雪花算法生成。对外我们暴露的是task_no,而不是自增 ID,防止用户遍历任务。索引设计上,(user_id, status)覆盖“用户查询自己任务列表”的场景,(status, created_at)覆盖“定时任务扫描超时任务”的场景。
3.2 异步消费与线程池隔离
任务进来之后,我用一个内存队列 + 消费者线程的方式执行。如果你的服务是多节点部署,应该把内存队列换成 Redis Stream 或者 RabbitMQ,思路完全一样。我目前的工作节点是单机部署,内存队列足够了,代码也更简单直接。
先定义任务执行器:
@Component public class TaskExecutor { private final ImageClient imageClient; private final ImageTaskMapper taskMapper; private final StringRedisTemplate redisTemplate; public void execute(ImageTask task) { updateStatus(task, "PROCESSING"); try { ImageData data = imageClient.generateWithMapping(task.getPrompt(), task.getSize(), 1); String finalUrl = persistImage(data, task.getTaskNo()); updateSuccess(task, finalUrl); } catch (ImageUpstreamException ex) { if (ex.getStatusCode() == 429) { // 触发限流,延迟重试 handleRateLimited(task); } else { handleRetryableFailure(task, ex); } } catch (Exception ex) { handleRetryableFailure(task, ex); } } }这里有一个关键决策:429 和 5xx 都走重试,但重试策略不一样。429 说明上游已经被打满,这时候如果你立刻重试只会雪上加霜,所以我用一个带退避的定时任务去延迟处理;5xx 则可能是瞬时故障,稍微退避一下就能恢复。
线程池这一层,考虑到我绝大多数时间阻塞在 HTTP 调用上,我直接用了虚拟线程:
@Bean("imageTaskExecutor") public ExecutorService imageTaskExecutor() { return Executors.newVirtualThreadPerTaskExecutor(); }如果不用虚拟线程,更稳妥的方案是定义一个固定大小的线程池,核心线程数 = CPU 核数 * 2,最大线程数 = 上游允许的并发数,队列容量根据任务提交速率设置。这里有一个经验值:线程池的并发上限不要超过上游 API Key 的 RPM(每分钟请求数)限制,否则你的重试机制会反过来成为放大请求的帮凶。
3.3 重试与退避策略:参数计算要落到实处
重试不是无脑循环,必须考虑三个参数:最大重试次数、重试间隔、退避因子。我推荐用指数退避加抖动(Exponential Backoff with Jitter)。纯固定间隔重试在某些场景下会形成“惊群效应”,所有失败请求同时重试,把上游再次打挂。抖动的作用就是让同一批次的重试请求在时间轴上分散开。
我用的是完整抖动策略,公式是:
sleep = random(0, min(cap, base * 2^attempt))代码实现:
public long nextBackoffMillis(int attempt, long baseMillis, long capMillis) { long exponent = Math.min(attempt, 10); // 防止指数溢出 long maxDelay = Math.min(capMillis, baseMillis * (long) Math.pow(2, exponent)); return ThreadLocalRandom.current().nextLong(0, maxDelay + 1); }以 base=1000ms、cap=30000ms 为例,各次重试前的随机休眠时间范围分别是:
| 重试次数 | 休眠范围 |
|---|---|
| 1 | 0 ~ 2s |
| 2 | 0 ~ 4s |
| 3 | 0 ~ 8s |
| 4 | 0 ~ 16s |
| 5 | 0 ~ 30s |
我通常把最大重试次数设为 3 到 5 次。超过这个次数,任务标记为 FAILED,然后走人工告警或者异步补偿。要注意的是,如果你的业务对图片时效性要求高,重试次数就少一点;如果允许延迟生成,可以放宽到 5 次。
重试的持久化这里也提一句:进程重启后,内存里的重试计划会全部丢失。如果你不能容忍这种情况,就需要把重试计划持久化到数据库或者 Redis ZSet,由定时任务扫描并重新投递。我因为是单机部署加上任务本身可以容忍少量失败,内存重试是够用的。
3.4 图片存储与 URL 有效期难题
gpt-image-2.5 返回的图片 URL 通常是有时效的,短则几十分钟,长则几天。如果你直接把上游 URL 存到数据库返回给用户,过段时间用户拿到的是 404。所以成熟做法是:拿到图片后立刻转存到自己的对象存储或者本地磁盘。
转存逻辑如下:
private String persistImage(ImageData data, String taskNo) { // 1. 如果返回的是 base64,直接写入磁盘 if (data.b64_json() != null) { byte[] imageBytes = Base64.getDecoder().decode(data.b64_json()); String filename = taskNo + ".png"; ossClient.putObject("ai-images", filename, new ByteArrayInputStream(imageBytes)); return buildCdnUrl(filename); } // 2. 如果返回的是 URL,先下载再转存 if (data.url() != null) { byte[] imageBytes = downloadFromUrl(data.url()); String filename = taskNo + ".png"; ossClient.putObject("ai-images", filename, new ByteArrayInputStream(imageBytes)); return buildCdnUrl(filename); } throw new IllegalStateException("unsupported image data format"); }下载上游图片这个过程,也需要设置合理的超时时间并限制图片大小。有些恶意响应可能给你回一个几个 GB 的流,直接把磁盘打满。我会用available()先判断 Content-Length,超过 20MB 的直接丢弃。
存储完之后,返回给用户的 URL 是你的 CDN 或 OSS 地址,就能保证长时间可访问了。这一节解决了“URL 过期”这个很多第一次接生图接口的人完全想不到的问题。
4. 防刷架构:从单点限流到多维风控
现在来到重头戏。防刷这件事,我单独开一节来讲,因为它的重要性和复杂度完全撑得起一个完整的章节。
4.1 防刷的本质:成本控制与公平使用
先算一笔账。假设一张图的上游成本是 0.1 元,你的服务日活用户是 1 万,如果完全不设防,恶意用户每人每天刷 1000 次,你的日成本是多少?100 万次请求,10 万元人民币。这还只是单一用户的攻击量,真实情况可能更夸张。
防刷的本质不是“防止所有用户使用”,而是在保证正常用户体验的前提下,识别并限制异常流量。它是一道成本闸门,不只是安全功能。
4.2 三层防线架构
我把防刷拆成三层,每层拦截不同维度的攻击:
第一层是网关层,用固定窗口或者令牌桶限制每个 IP 的请求频率。这个层级的拦截目标是脚本和爬虫,它们在短时间内发起大量同 IP 请求,特征非常明显。这一层我用 Spring Boot 的拦截器实现,基于 Redis 计数。
第二层是业务层,核心是用户维度限流和额度控制。用户登没登录、有没有剩余配额、单位时间内允许多少次调用,都在这一层校验。这一层拦截的是“已登录但滥用”的用户,他们可能用合法账号做违规操作。
第三层是数据层,通过幂等键和历史行为分析防止重放攻击。比如同一个请求 ID 短时间出现多次,直接拒绝;某个用户的历史请求量异常突增,触发风险控制规则。
三层设计避免了单点失败:如果第一层被绕过,第二层还能挡住;第二层有疏漏,第三层还能补救。你可以根据自身业务量级决定上几层,但最少要有两层。
4.3 核心实现:令牌桶限流 + 用户额度 + 幂等控制
令牌桶限流我基于 Redis 的 Lua 脚本实现。为什么不直接INCR加EXPIRE?因为固定窗口限流存在临界问题:用户在窗口末尾和下一个窗口开始连续请求,实际速率可以达到限制的两倍。令牌桶算法能平滑突发流量,这也是生产环境更倾向于令牌桶的原因。
我用 Redisson 的RRateLimiter做令牌桶,它底层也是 Lua 脚本,API 封装得不错:
@Configuration public class RateLimitConfig { @Bean public RRateLimiter ipRateLimiter(RedissonClient redisson) { RRateLimiter limiter = redisson.getRateLimiter("ip:limiter"); limiter.trySetRate(RateType.OVERALL, 10, 1, RateIntervalUnit.MINUTES); return limiter; } }这里10表示令牌桶容量,1表示每分钟补充 10 个令牌。也就是说,允许每 6 秒一个请求,但允许瞬间突发 10 个。具体数值要根据你的业务调整,工具类网站的访问密度和内部系统的访问密度完全不是一个量级。
拦截器中使用:
@Component public class RateLimitInterceptor implements HandlerInterceptor { private final RRateLimiter limiter; @Override public boolean preHandle(HttpServletRequest request, HttpServletResponse response, Object handler) throws Exception { String clientIp = request.getRemoteAddr(); String key = "rate:ip:" + clientIp; RRateLimiter userLimiter = redisson.getRateLimiter(key); userLimiter.trySetRate(RateType.OVERALL, 10, 1, RateIntervalUnit.MINUTES); if (!userLimiter.tryAcquire(1)) { response.setStatus(429); response.getWriter().write("{\"code\":429,\"message\":\"too many requests\"}"); return false; } return true; } }用户额度控制是另一道独立的逻辑。每个用户有一个每日配额,比如普通用户 50 张、会员用户 500 张。实现上我用 Redis 的 Hash 结构存储用户当天已用次数,键名带日期,比如quota:user:1001:2025-06-17,字段为used。每次请求前先检查是否达到上限,请求成功后执行INCR。这个逻辑放在 Service 层,用 AOP 注解会很优雅,但为了可读性我这里直接写在方法里:
public String submitTask(ImageGenerateDTO dto, Long userId) { // 额度检查 String quotaKey = "quota:user:" + userId + ":" + LocalDate.now(); String used = redisTemplate.opsForValue().get(quotaKey); int usedCount = used == null ? 0 : Integer.parseInt(used); if (usedCount >= getUserDailyLimit(userId)) { throw new QuotaExceededException("daily quota exceeded"); } // 创建任务(任务表记录 userId) ImageTask task = createTask(dto, userId); // 入队 taskQueue.submit(task); // 扣减额度(也可以等任务成功后再扣,这里说明一下差异) redisTemplate.opsForValue().increment(quotaKey, 1); redisTemplate.expire(quotaKey, Duration.ofHours(26)); return task.getTaskNo(); }关于扣减时机,我多说一句:按提交时扣适合“先付费后使用”的模式,抗恶意提交但用户可能遇到“提交了但失败,额度白扣”;按成功时扣体验更友好,但恶意用户会把任务全部提交进队列,让系统资源被无效任务占满。我的做法是提交时扣,但失败后自动返还(调用失败回调里执行decrement)。这样既防了资源占用,又不损害正常用户体验。
幂等控制这块很容易被忽略,但它其实是防刷里性价比最高的手段。用户或攻击者重复提交同一个请求,如果每次都能创建新任务,那即使有额度限制,也能通过大量重复请求耗光你的队列资源。幂等控制的做法很简单:客户端提交请求时携带一个Idempotency-Key,后端在 Redis 里查这个 Key 是否已存在,存在就直接返回第一次的任务号,否则创建新任务并缓存 Key。
String idempotencyKey = request.getHeader("Idempotency-Key"); if (idempotencyKey != null) { Boolean firstTime = redisTemplate.opsForValue() .setIfAbsent("idem:" + idempotencyKey, taskNo, Duration.ofHours(1)); if (!firstTime) { String existingTaskNo = redisTemplate.opsForValue().get("idem:" + idempotencyKey); return existingTaskNo; } }4.4 阈值参数怎么定:从数据和成本反推
限流阈值不是拍脑袋定的。我提供一个反推方法。
假设你的上游接口配额是每分钟 600 张(RPM=600),你的服务有 2 个节点,那你所有节点的总速率应该控制在 500 左右,给上游留 20% 余量。每个节点就是 250 张/分钟,换算成 QPS 是 4.2。你想 90% 的配额给正常用户,10% 给突发缓冲,那正常用户的 IP 级限流大约就是每分钟 4 次左右。注意,这是全局限流,不是单 IP 限流。
单 IP 限流应该参考“真人用户的操作频率”。一个正常人看 AI 生图网站,每分钟提交 1 到 2 次任务已经是相当高频了。所以单 IP 每分钟 10 次是一个合理的起点,如果业务有明显的高峰操作需求,可以放开到每分钟 20 次,超过这个量基本可以断定不是真人。
用户额度上限参考你的成本预算。你每天预算 1000 元,单张成本 0.1 元,那全天的生成上限就是 10000 张。如果注册用户是 200 人,人均 50 张就能覆盖预算。这些数字都要落到配置中心,方便随时调整,不要写死在代码里。
4.5 防刷策略下的用户体验权衡
防刷做太狠会误伤正常用户,这也是我踩过坑的地方。比如某次我把单 IP 限流调到每分钟 3 次,公司内部测试群里立刻有人反馈“图片转半天出不来”,一看日志全是 429。所以现在我的策略是:IP 限流只做粗粒度的脚本拦截,精确控制靠用户维度的配额和风控规则。IP 限流阈值可以放宽到正常用户操作的 5 到 10 倍;用户维度配额严格限制;再追加一个风险规则:如果同一用户短时间提交超过 N 次,触发人机验证。
5. 完整流程串讲与关键代码走读
架构和核心模块都过完了,我把整个流程串起来走一遍,从输入请求到返回图片 URL 的完整链路。
5.1 控制层与参数校验
控制层的职责非常薄,只做参数格式校验和调用 Service。我再三强调一个原则:Controller 里不要写业务逻辑。所有情况都让 Service 层处理,Controller 只负责 HTTP 语义的适配。
@RestController @RequestMapping("/api/v1/images") public class ImageController { private final ImageTaskService taskService; @PostMapping("/generations") public ResponseEntity<SubmitResponse> submit( @Valid @RequestBody ImageGenerateDTO dto, @RequestHeader(value = "Idempotency-Key", required = false) String idempotencyKey, @RequestAttribute("userId") Long userId) { String taskNo = taskService.submit(dto, userId, idempotencyKey); return ResponseEntity.accepted() .body(new SubmitResponse(taskNo, "task accepted")); } @GetMapping("/tasks/{taskNo}") public ResponseEntity<TaskResult> query(@PathVariable String taskNo) { TaskResult result = taskService.query(taskNo); return ResponseEntity.ok(result); } }@RequestAttribute("userId")是从拦截器里解析出来的登录态用户 ID。这个项目里我用 JWT 做认证,拦截器解析 Token 并注入用户信息,因为篇幅关系就不展开 JWT 的前后端交互细节了。
参数校验这块,我写了一个简单的 DTO 校验规则:
public record ImageGenerateDTO( @NotBlank(message = "prompt cannot be empty") @Size(max = 2000, message = "prompt too long") String prompt, @Pattern(regexp = "512x512|1024x1024", message = "unsupported size") String size ) {}prompt最长 2000 个字符是常规限制,防止有人提交超长文本把上游接口打异常;size只允许后端配置的值,避免用户传任意值,这个也是成本和稳定性考虑。
5.2 任务提交与异步执行串讲
Service 层的submit方法把额度检查、幂等判断、任务创建、入队四件事串起来。前面已经贴过这一段的代码逻辑,这里补一下任务入队的执行代码:
@Component public class TaskQueue { private final ExecutorService executor; private final TaskExecutor taskExecutor; public void submit(ImageTask task) { executor.submit(() -> { try { taskExecutor.execute(task); } catch (Exception e) { log.error("task execute failed. taskNo={}", task.getTaskNo(), e); } }); } }这里executor.submit返回的Future我没有接收,因为异常已经在taskExecutor.execute内部处理过并映射到任务状态了。如果这里再包一层 Future 等待,就会阻塞提交线程,失去异步的意义。
查询接口也一起说一下:
public TaskResult query(String taskNo) { ImageTask task = taskMapper.selectByTaskNo(taskNo); if (task == null) { throw new TaskNotFoundException("task not found"); } return new TaskResult( task.getTaskNo(), task.getStatus(), task.getImageUrl(), task.getFailureReason() ); }前端拿到任务状态后,如果是PROCESSING就轮询这个查询接口,每隔 2 到 3 秒查一次;如果是SUCCEEDED就展示图片;如果是FAILED就提示用户重新提交。这个轮询间隔也要做说明:2 到 3 秒是比较平衡的节奏,太短会对服务造成无谓压力,太长又影响用户体验。你也可以换成 WebSocket 推送或者 SSE,但那是架构升级,不是必需的。
5.3 定时任务兜底:扫描超时与失败补偿
异步任务流程中,最怕出现“任务卡死在某一个中间状态”。比如消费者线程拿到了任务,正在调上游接口时进程崩溃,任务状态停留在PROCESSING永远不更新。所以我加了一个定时补偿任务,每 30 秒扫描一次数据库,将所有超过 5 分钟仍然处于PROCESSING或QUEUED状态的任务找出来,统一标记为TIMEOUT。
@Scheduled(fixedDelay = 30000) public void scanTimeoutTasks() { List<ImageTask> stuckTasks = taskMapper.selectStuckTasks( "PROCESSING", LocalDateTime.now().minusMinutes(5) ); for (ImageTask task : stuckTasks) { log.warn("stuck task detected. taskNo={}, status={}", task.getTaskNo(), task.getStatus()); taskMapper.updateStatus(task.getId(), "TIMEOUT", task.getVersion()); } }这个兜底逻辑保证了任务不会无限期挂起。如果你用了 Redis Stream 或 MQ,还需要处理消费者宕机后的消息重新入队,这里用数据库状态做兜底就足够简单可靠了。
6. 常见问题与排查技巧实录
这一节的内容全部来自我真实操作中踩过的坑,写成速查表供你排查问题的时候对照使用。
6.1 高频问题速查表
| 现象 | 可能原因 | 排查方式 | 解决方案 |
|---|---|---|---|
| 上游返回 401 | API Key 错误或过期 | 检查请求 Header 中的 Authorization | 从配置中心更新 Key,注意别缓存到内存 |
| 上游返回 429 | 触发上游限流 | 查看响应头和日志时间戳 | 降低并发,调整重试退避策略 |
| 任务一直 PROCESSING | 消费者线程池耗尽 | 查看线程池活跃线程数 | 调大线程池上限或者查上游 RT 是否飙高 |
| 图片 URL 访问 404 | 上游 URL 过期未转存 | 检查任务表 image_url 字段前缀 | 确认转存逻辑执行成功,OSS 权限正确 |
| 偶发 5xx 重试后成功 | 上游瞬时故障 | 看日志中错误码是否符合重试条件 | 确保 5xx 和超时都纳入重试,4xx 不重试 |
| 用户额度莫名减少 | 提交时扣减策略导致失败也扣 | 核对扣减与返还日志 | 失败自动返还额度 |
| Redis 限流计数不准 | 多个节点各计各的 | 检查 Redis 键的命名规则 | 确认 key 中带用户 ID 或 IP,用统一 Redis |
| 请求被网关层拦截误伤内测用户 | IP 限流阈值太低 | 看被拦截的 IP 和时间分布 | 提高阈值或加白名单 |
6.2 好几个真实踩坑记录
第一个坑:虚拟线程与连接池的配合问题。我一开始用虚拟线程 + JdkClientHttpRequestFactory,表现非常好,1 万个并发不费劲。但后来发现,上游接口偶尔会挂起很长时间,导致虚拟线程大量堆积在等待响应上,本地连接数直接冲到几十万。排查之后才明白:虚拟线程虽好,但如果你没设置读超时,等待时间理论上可以无限长。后来我在 RestClient 那里加了读取超时,情况立刻缓解。所以再次提醒:虚拟线程释放的是平台线程,不是释放连接资源,超时设置永远是第一位的。
第二个坑:幂等 Key 设置了太多过期时间。某个版本我把幂等 Key 的过期时间设成了 7 天,结果用户同一个 Key 在一周内提交,永远拿到第一次的任务号,导致后续任务完全无法生成。后来想明白,幂等窗口只需要覆盖“客户端可能重试的最大时间区间”就行了,正常设 1 小时足够。时间太短客户端重试会重复提交,太长又影响正常使用。这是设计上的细节,容易忽略。
第三个坑:积分扣减与任务成功的时序问题。我之前用的方案是“任务成功后再扣减额度”,结果被刷过一次。攻击者用低并发、多账号的方式绕过了 IP 限流,把所有任务都提交进队列,排队等待生成,但队列被占满后正常用户的请求全部超时。换成提交时扣减后,攻击者虽然还能提交任务,但他的额度很快用完,对正常用户的影响就小多了。这也是生产事故教会我的一个经验。
第四个坑:重试风暴。有一次上游接口出现故障,所有请求都返回 5xx。我的重试逻辑一看是 5xx 就立刻重新入队,结果形成了“失败-重试-失败-重试”的死循环,上游请求量直接翻了几倍。现在我的重试策略里加了一个熔断开关:如果连续 30 秒内错误率超过 50%,直接打开熔断器,所有生图任务快速失败,5 分钟后再尝试恢复。这个熔断逻辑虽然简单粗暴,但确实能保证系统不被拖垮。
6.3 监控与告警配置
最后说一下监控。生图服务有几个核心指标我每天都看:
- 任务提交速率:判断流量是否异常。
- 任务成功率:正常应该在 95% 以上,低于这个值要立刻查日志。
- 上游接口 RT 的 P50/P95/P99:P95 超过 30 秒就要考虑是否要调整读取超时和重试策略。
- 队列积压数量:内存队列无界会导致内存增长,所以我在提交入口加了队列长度检查,超过 10000 直接返回“系统繁忙”。
- 用户额度使用率:每天跑一个报表,看是否有用户接近配额上限,提前预判成本。
Actuator 暴露/actuator/metrics和/actuator/health,Prometheus 每 15 秒抓一次,配合 Grafana 做面板。告警规则只设两条最关键的:成功率低于 90% 持续 5 分钟告警;队列积压连续 3 分钟超过 5000 告警。告警太多等于没有告警,把有限的精力集中在最重要的事情上。
写在最后
这套管道和防刷架构,我在几个实际项目里跑过,整体稳定性还是不错的。最关键的一个体会是:Go for production,要先把“失败”想清楚。什么时候会失败、失败之后怎么恢复、恢复不了怎么兜底,这些问题比“怎么调通接口”重要得多。设计阶段多想一层,运维阶段就少熬一夜。
后面如果你想扩展,可以考虑把内存队列换成 Redis Stream 做多节点水平扩展,或者引入 MQ 做削峰填谷。但如果你的业务量级就是单个节点能撑住的量,这套架构已经足够你撑过业务从 0 到 1 的阶段了。