ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Java AI应用异步化与高并发实战:从虚拟线程到WebFlux

Java AI应用异步化与高并发实战:从虚拟线程到WebFlux 1. 为什么Java AI应用必须直面异步化与高并发这道坎最近三个月我连续接手了三个AI类后端项目一个实时多模态内容审核服务、一个面向金融风控的实时推理API网关、还有一个教育领域的个性化学习路径生成引擎。它们表面看都是“调用大模型API”但上线后无一例外在压测阶段暴露出同一个致命问题——吞吐量卡死在每秒30~50请求CPU利用率常年95%以上GC频繁到日志里全是Full GC警告。排查下来核心症结不是模型本身慢而是Java线程池被阻塞得像早高峰地铁站——每个HTTP请求进来就占住一个Tomcat线程而这个线程要等外部AI服务返回结果才能释放。更糟的是这些AI服务响应时间极不稳定有时200ms有时3秒甚至偶尔超时重试。线程池里积压的请求越来越多新请求进不来老请求出不去整个系统进入雪崩前夜。这就是典型的同步阻塞式AI调用陷阱。很多人以为“Java Spring Boot OpenAI SDK”就能跑通AI应用却忽略了AI服务天然具备的三大非理想特性网络延迟不可控、响应时间长尾分布严重、失败率远高于传统微服务。拿我们那个内容审核服务举例它需要同时调用图像识别、文本情感分析、OCR三个AI接口每个接口平均耗时800msP95达到2.3秒。如果用传统RestTemplate同步调用一个请求就要独占线程近3秒而Tomcat默认线程池只有200个线程——理论上最大QPS才66实际还要打七折。这根本撑不住真实业务场景下每秒数百甚至上千的并发请求。真正破局的关键在于把“等待AI响应”这件事从主线程中剥离出去。不是简单加个Async注解就完事而是要构建一套完整的异步化流水线请求进来立刻返回202 Accepted后台用独立线程池处理AI调用结果通过WebSocket或消息队列推送给前端中间状态可查询、可取消、可重试。这套设计背后是Java生态里几代人沉淀下来的高并发智慧——从Future到CompletableFuture从Reactor到WebFlux再到Spring Boot 3.x全面拥抱虚拟线程Virtual Threads。我亲眼见过一个团队把旧版同步服务迁移到基于虚拟线程的异步架构后同样硬件资源下QPS从42飙升到1800线程数从200降到不到50个。这不是玄学而是把Java最擅长的“调度”能力精准用在了AI应用最脆弱的环节上。你可能正在面试Java后端岗刷到“Java多线程和高并发”这类题也可能正开发一个AI Agent系统被用户投诉“响应太慢”。无论哪种情况理解异步化与高并发的设计逻辑已经不再是加分项而是生存必需。它不等于堆砌技术名词而是要搞懂当一个HTTP请求触发一次LLM调用时你的代码里到底发生了什么线程在哪创建在哪阻塞在哪释放数据如何在线程间安全传递错误如何优雅降级接下来我会拆解整套方案从原理到代码从选型到避坑全部基于真实生产环境验证过的方法。2. 异步化与高并发设计的整体架构与选型逻辑2.1 为什么不能只靠Async——同步阻塞的底层真相很多开发者第一反应是给AI调用方法加上Async注解以为这就实现了异步。但我在三个项目里都踩过这个坑表面看请求不卡了实际系统负载反而更高。原因在于Async只是把方法扔进另一个线程池执行它并没有解决I/O阻塞的本质问题。举个具体例子Service public class AiService { Async // 这行代码看似解决了问题 public void analyzeImage(String imageUrl) { // 这里调用RestTemplate.getForObject(...) // 线程在此处被网络I/O阻塞无法做其他事 String result restTemplate.getForObject(https://ai-api.com/v1/analyze, String.class, imageUrl); saveResult(result); } }这段代码的问题在于restTemplate.getForObject()内部使用的是阻塞式HTTP客户端如Apache HttpClient它会一直占用当前线程直到响应返回。即使这个方法在独立线程池里执行该线程在等待网络响应时也是“挂起”状态无法被调度器复用。假设你配置了20个Async线程同时有50个图片分析请求进来其中30个请求卡在网络等待上剩下20个线程根本无法启动新任务——线程池被无效占满系统吞吐量不升反降。真正的异步必须是非阻塞I/ONon-blocking I/O。这意味着HTTP客户端在发起请求后立即返回一个Future或Mono对象主线程可以继续处理其他任务当网络响应到达时由事件循环Event Loop或回调机制通知业务逻辑处理结果。Java生态里实现这一点的主流方案有两条技术路线Reactive路线基于Netty的WebClientSpring WebFlux、Project Reactor的Mono/Flux完全异步非阻塞内存占用低适合高并发长连接场景。Virtual Threads路线JDK 21的虚拟线程Project Loom用类似协程的方式让阻塞式代码也能高效运行开发体验接近传统编程迁移成本低。我对比过两种方案在AI服务中的表现。用WebClient重构后单机QPS提升3.2倍但开发复杂度显著增加——所有DAO层都要改成响应式如R2DBC异常处理链路变长调试难度上升。而用虚拟线程方案只需把RestTemplate换成HttpClient支持虚拟线程再把Controller方法标记为Transactional即可代码改动量不到20%QPS提升2.8倍。对于大多数已有Spring Boot项目我强烈推荐虚拟线程路线作为首选它把“异步”的门槛降到了最低。2.2 高并发下的流量整形与熔断策略——别让AI服务拖垮整个系统AI服务不是数据库它没有连接池概念也没有事务一致性要求但它有更严苛的限制速率限制Rate Limiting和配额Quota。OpenAI官方API明确要求每分钟最多60次请求免费 tierAzure AI服务按Token数计费国内大模型平台普遍采用QPS并发数双重限制。如果放任前端直接调用瞬间涌来的1000个请求会全部打到AI服务上轻则触发限流返回429重则导致账号被封禁。因此高并发设计的第一道防线是服务端流量整形。我不会用简单的RateLimiter注解而是采用分层控制接入层限流Nginx配置limit_req按IP或API Key限制QPS防止单点恶意请求打穿网关网关层熔断Spring Cloud Gateway集成Resilience4j对AI服务调用设置failureRateThreshold50%连续5次失败自动熔断30秒业务层排队用BlockingQueue实现公平队列所有AI请求先进队列再由固定数量工作线程如10个按顺序消费确保不超出AI服务商的QPS上限。这里有个关键细节队列长度不能无限大。我见过一个项目把队列设为Integer.MAX_VALUE结果突发流量导致内存溢出OOM。正确做法是根据SLA计算最大等待时间。比如AI服务P95响应时间2秒业务要求用户等待不超过10秒那么队列容量10秒 / 2秒 * 工作线程数5*1050。超过50的请求直接返回429 Too Many Requests并附带Retry-After: 5头提示客户端5秒后重试。第二道防线是智能降级。当AI服务不可用时不能简单返回错误。我们设计了三级降级策略一级返回缓存的历史相似结果如用户上次提问的答案二级调用轻量级本地模型如TinyBERT生成基础回复三级返回预设的兜底话术如“AI正在思考请稍候”。这套策略在金融风控项目中救了我们一命——某天OpenAI服务大面积超时系统自动切换到本地模型虽然准确率下降15%但业务连续性100%保障客户投诉率为零。2.3 数据一致性与状态管理——异步流程中的“事务”难题异步化最大的挑战不是技术实现而是状态一致性。同步调用时一个HTTP请求对应一个数据库事务成功则提交失败则回滚。但异步流程中请求接收、AI调用、结果存储、通知推送可能分布在不同线程甚至不同服务中。如何保证“用户提交审核请求→AI分析完成→结果入库→前端收到通知”这一串操作要么全成功要么全失败我的方案是Saga模式 状态机。不依赖分布式事务XA协议在高并发下性能极差而是把长流程拆成一系列本地事务每个步骤都有对应的补偿操作。以内容审核为例步骤主操作补偿操作触发条件1创建审核任务记录statusPROCESSING删除任务记录步骤1失败2调用AI服务获取结果调用AI服务取消任务步骤2超时或失败3更新任务状态为SUCCESS/FAILED恢复任务状态为PROCESSING步骤3失败状态机用StateMachine框架如Spring Statemachine管理所有状态变更都通过事件驱动。关键点在于补偿操作必须幂等。比如“取消AI任务”接口要支持重复调用避免因网络重试导致误取消。我们给每个AI请求生成唯一requestId所有操作都带上这个ID服务端用Redis的SETNX保证幂等性。另外用户需要实时感知进度。我们不用轮询浪费资源而是用Server-Sent EventsSSE推送状态更新。Controller返回SseEmitterAI调用完成后通过emitter.send()推送JSON事件前端用EventSource监听。实测在500并发下SSE比WebSocket内存占用低40%且兼容性更好无需额外WebSocket服务器。3. 核心实现细节与实操要点3.1 基于虚拟线程的异步AI调用——零改造迁移方案虚拟线程Virtual Threads是JDK 21的里程碑特性它让Java终于拥有了类似Go协程的轻量级线程。一个虚拟线程仅占用几KB栈空间而传统平台线程Platform Thread动辄1MB。这意味着你可以轻松创建10万个虚拟线程处理10万并发请求而平台线程池撑死也就几百个。要启用虚拟线程第一步是升级JDK。我们线上环境已全面切换到JDK 21Spring Boot版本需≥3.2Spring Framework 6.1原生支持虚拟线程。关键配置如下# application.yml spring: threads: virtual: enabled: true # 启用虚拟线程支持 web: server: tomcat: max-connections: 10000 # 提高连接数 accept-count: 1000 # 提高等待队列长度第二步替换阻塞式HTTP客户端。RestTemplate不支持虚拟线程必须改用java.net.http.HttpClientJDK 11内置Configuration public class HttpClientConfig { Bean public HttpClient httpClient() { // 关键启用虚拟线程调度器 return HttpClient.newBuilder() .executor(Executors.newVirtualThreadPerTaskExecutor()) .connectTimeout(Duration.ofSeconds(10)) .build(); } }第三步改造AI服务调用。注意不要用Future.get()阻塞等待这是虚拟线程的大忌。正确写法是用CompletableFuture配合thenApply链式处理Service public class AiAnalysisService { private final HttpClient httpClient; public AiAnalysisService(HttpClient httpClient) { this.httpClient httpClient; } // 返回CompletableFuture不阻塞主线程 public CompletableFutureString analyzeText(String text) { HttpRequest request HttpRequest.newBuilder() .uri(URI.create(https://ai-api.com/v1/text-analyze)) .header(Content-Type, application/json) .POST(HttpRequest.BodyPublishers.ofString({\text\:\ text \})) .build(); // 异步发送请求返回CompletableFuture return httpClient.sendAsync(request, HttpResponse.BodyHandlers.ofString()) .thenApply(HttpResponse::body) .exceptionally(throwable - { log.error(AI分析失败, throwable); return {\error\:\timeout\}; }); } }Controller层直接返回CompletableFutureSpring Boot会自动处理RestController public class AiController { private final AiAnalysisService aiService; public AiController(AiAnalysisService aiService) { this.aiService aiService; } PostMapping(/api/analyze) public CompletableFutureResponseEntityString analyze(RequestBody String text) { return aiService.analyzeText(text) .thenApply(result - ResponseEntity.ok(result)) .exceptionally(throwable - ResponseEntity.status(500).build()); } }实测效果同一台4核8G服务器同步方案最大QPS 42虚拟线程方案达1180线程数从217降至38。更重要的是代码几乎没变——只是把RestTemplate换成HttpClient把return改成return CompletableFuture开发成本极低。提示虚拟线程不是银弹。它不适合CPU密集型任务如本地模型推理因为虚拟线程调度器会把CPU密集任务交给平台线程执行。AI调用本质是I/O密集型所以完美匹配。3.2 响应式WebFlux方案——极致性能的终极选择当业务对延迟极度敏感如实时聊天机器人或者需要处理海量长连接如万级WebSocket连接WebFlux是更优解。它基于Netty事件驱动所有操作都在一个线程池EventLoopGroup内完成避免线程上下文切换开销。第一步创建WebFlux项目。Spring Initializr选择Spring Reactive Web依赖自动包含spring-boot-starter-webflux。关键配置# application.yml spring: web: flux: response-timeout: 30s # 全局响应超时 redis: host: localhost port: 6379第二步用WebClient替代RestTemplate。WebClient是响应式HTTP客户端返回Mono或FluxService public class ReactiveAiService { private final WebClient webClient; public ReactiveAiService(WebClient.Builder webClientBuilder) { this.webClient webClientBuilder .codecs(configurer - configurer.defaultCodecs().maxInMemorySize(10 * 1024 * 1024)) // 提高内存限制 .build(); } public MonoString analyzeImage(String imageUrl) { return webClient.post() .uri(https://ai-api.com/v1/image-analyze) .header(Content-Type, application/json) .bodyValue({\url\:\ imageUrl \}) .retrieve() .bodyToMono(String.class) .timeout(Duration.ofSeconds(15)) // 设置超时 .onErrorResume(e - { log.error(AI图像分析失败, e); return Mono.just({\error\:\service_unavailable\}); }); } }第三步Controller返回Mono并集成Redis做结果缓存RestController public class ReactiveAiController { private final ReactiveAiService aiService; private final RedisTemplateString, Object redisTemplate; public ReactiveAiController(ReactiveAiService aiService, RedisTemplateString, Object redisTemplate) { this.aiService aiService; this.redisTemplate redisTemplate; } PostMapping(/api/analyze-reactive) public MonoResponseEntity? analyzeReactive(RequestBody String payload) { String cacheKey ai_result: DigestUtils.md5Hex(payload); // 先查缓存 return redisTemplate.opsForValue().get(cacheKey) .cast(String.class) .flatMap(cached - Mono.just(ResponseEntity.ok(cached))) .switchIfEmpty( // 缓存未命中调用AI服务 aiService.analyzeImage(payload) .flatMap(result - { // 结果存入Redis过期时间30分钟 redisTemplate.opsForValue().set(cacheKey, result, Duration.ofMinutes(30)); return Mono.just(ResponseEntity.ok(result)); }) .onErrorResume(e - Mono.just(ResponseEntity.status(500).build())) ); } }性能对比WebFlux方案在万级并发下P99延迟稳定在80ms以内而Servlet容器方案P99达1200ms。但代价是开发心智负担重——所有DAO层必须用R2DBC响应式数据库驱动不能混用JPA/Hibernate。我们只在核心AI网关服务中采用此方案普通业务服务仍用虚拟线程。3.3 高并发下的线程池与资源隔离——避免“雪崩式传染”异步化后线程池管理变得比以往更重要。一个设计不良的线程池会让整个系统陷入瘫痪。我总结出三条铁律第一绝不共享线程池。很多项目用一个Async线程池处理所有异步任务结果AI调用慢导致线程池占满连邮件发送、日志上传都卡住。正确做法是按任务类型划分线程池Configuration public class ThreadPoolConfig { // AI调用专用线程池核心线程数CPU核数最大线程数20队列容量100 Bean(aiThreadPool) public Executor aiThreadPool() { ThreadPoolTaskExecutor executor new ThreadPoolTaskExecutor(); executor.setCorePoolSize(Runtime.getRuntime().availableProcessors()); executor.setMaxPoolSize(20); executor.setQueueCapacity(100); executor.setThreadNamePrefix(ai-pool-); executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy()); // 拒绝策略由调用线程执行 executor.initialize(); return executor; } // 邮件发送线程池核心2最大5队列10邮件发送本身不耗时 Bean(emailThreadPool) public Executor emailThreadPool() { // ... 类似配置 return executor; } }第二合理设置队列容量。无界队列LinkedBlockingQueue无参构造是毒药。我们曾用无界队列结果OOM崩溃。现在一律用有界队列并设置拒绝策略为CallerRunsPolicy——当队列满时由调用方线程自己执行任务。这能自然形成背压Backpressure让上游慢下来避免雪崩。第三监控线程池健康度。Spring Boot Actuator提供/actuator/threaddump端点但我们更需要实时指标。我用Micrometer集成PrometheusBean public MeterRegistryCustomizerMeterRegistry metrics() { return registry - { // 监控AI线程池 ThreadPoolTaskExecutor aiPool (ThreadPoolTaskExecutor) applicationContext.getBean(aiThreadPool); registry.gauge(threadpool.ai.active, aiPool, pool - (double) pool.getActiveCount()); registry.gauge(threadpool.ai.queue.size, aiPool, pool - (double) pool.getThreadPoolExecutor().getQueue().size()); }; }Grafana面板实时显示threadpool.ai.queue.size 80就告警运维人员立即扩容或检查AI服务状态。3.4 状态持久化与结果推送——异步流程的闭环设计异步流程的终点不是“调用AI成功”而是“用户看到结果”。这需要一套可靠的状态持久化与推送机制。我们的方案是数据库状态表 Redis缓存 SSE推送三者结合。首先设计状态表ai_taskCREATE TABLE ai_task ( id BIGINT PRIMARY KEY AUTO_INCREMENT, task_id VARCHAR(64) NOT NULL COMMENT 全局唯一任务ID, user_id BIGINT NOT NULL COMMENT 用户ID, status ENUM(PENDING, PROCESSING, SUCCESS, FAILED, CANCELLED) DEFAULT PENDING, request_data TEXT COMMENT 原始请求数据, result_data TEXT COMMENT AI返回结果, created_time DATETIME DEFAULT CURRENT_TIMESTAMP, updated_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, INDEX idx_user_status (user_id, status), INDEX idx_task_id (task_id) );关键点task_id用雪花算法生成SnowflakeIdGenerator确保全局唯一且有序status字段用ENUM类型避免字符串拼写错误添加复合索引加速按用户查询任务列表。其次用Redis缓存高频访问的任务状态。我们定义缓存Key为ai:task:{taskId}TTL设为24小时业务要求结果保留一天Service public class TaskStatusService { private final RedisTemplateString, Object redisTemplate; private final JdbcTemplate jdbcTemplate; public void updateStatus(String taskId, String status, String result) { // 1. 更新数据库 String sql UPDATE ai_task SET status ?, result_data ?, updated_time NOW() WHERE task_id ?; jdbcTemplate.update(sql, status, result, taskId); // 2. 更新Redis缓存 String cacheKey ai:task: taskId; MapString, Object cacheData new HashMap(); cacheData.put(status, status); cacheData.put(result, result); cacheData.put(updatedAt, System.currentTimeMillis()); redisTemplate.opsForHash().putAll(cacheKey, cacheData); redisTemplate.expire(cacheKey, Duration.ofHours(24)); } }最后SSE推送。Controller中创建SseEmitter并存入内存Map生产环境建议用Redis Pub/SubRestController public class TaskStatusController { // 内存Map存储Emitter生产环境应替换为Redis private final MapString, SseEmitter emitterMap new ConcurrentHashMap(); GetMapping(/api/task/{taskId}/status) public SseEmitter getStatus(PathVariable String taskId) { SseEmitter emitter new SseEmitter(30 * 60 * 1000L); // 30分钟超时 // 任务完成时推送 emitter.onCompletion(() - emitterMap.remove(taskId)); emitter.onError(throwable - { log.error(SSE连接异常, throwable); emitterMap.remove(taskId); }); emitterMap.put(taskId, emitter); return emitter; } // AI服务回调时触发推送 Service public class AiCallbackService { public void onTaskComplete(String taskId, String result) { SseEmitter emitter emitterMap.get(taskId); if (emitter ! null) { try { emitter.send(SseEmitter.event() .name(status-update) .data({\status\:\SUCCESS\,\result\: result })); emitter.complete(); emitterMap.remove(taskId); } catch (IOException e) { log.error(推送SSE失败, e); emitterMap.remove(taskId); } } } } }前端JavaScript监听const eventSource new EventSource(/api/task/abc123/status); eventSource.addEventListener(status-update, event { const data JSON.parse(event.data); if (data.status SUCCESS) { document.getElementById(result).innerText data.result; } });这套方案实测在2000并发SSE连接下内存占用稳定在1.2GBCPU利用率低于30%。比轮询方案节省90%的带宽和服务器资源。4. 常见问题与排查技巧实录4.1 “明明用了Async为什么还是卡”——线程池陷阱全解析这是最高频的问题。我整理了五种典型场景及解决方案场景现象根本原因解决方案场景1未配置自定义线程池所有Async方法共用SimpleAsyncTaskExecutor每次新建线程OOM风险高EnableAsync默认使用SimpleAsyncTaskExecutor不复用线程显式配置ThreadPoolTaskExecutorBean设置核心/最大线程数场景2事务失效Async方法内数据库操作不生效或事务不回滚Async方法在新线程执行脱离原事务上下文在Async方法内手动开启事务TransactionTemplate或改用Transactional(propagation Propagation.REQUIRES_NEW)场景3异常丢失Async方法抛异常调用方完全不知情Async默认不传播异常异常被吞掉在线程池配置中设置setWaitForTasksToCompleteOnShutdown(true)或用Future.get()捕获异常场景4循环依赖Service A调用B的Async方法B又注入A启动报错Spring代理机制导致循环引用改用ApplicationContext.getBean()获取Bean或重构为事件驱动场景5静态方法无效给static方法加Async完全不生效Async基于Spring AOP代理静态方法无法被代理去掉static改为实例方法最隐蔽的陷阱是场景3异常丢失。我曾遇到一个AI服务因网络超时抛出SocketTimeoutException但Controller层收不到任何错误用户看到空白页面。排查方法是在ThreadPoolTaskExecutor中添加异常处理器Bean(aiThreadPool) public Executor aiThreadPool() { ThreadPoolTaskExecutor executor new ThreadPoolTaskExecutor(); // ... 其他配置 executor.setRejectedExecutionHandler((r, executor1) - { log.error(AI线程池拒绝任务, new RuntimeException(Rejected)); }); executor.setThreadFactory(r - { Thread t new Thread(r); t.setUncaughtExceptionHandler((thread, ex) - { log.error(AI线程未捕获异常, ex); }); return t; }); return executor; }这样所有未捕获异常都会记录到日志便于快速定位。4.2 “AI服务响应忽快忽慢怎么优化”——网络与序列化瓶颈定位AI服务响应时间波动大往往不是AI模型问题而是客户端侧的网络或序列化瓶颈。我用三步法定位第一步排除DNS解析延迟。很多团队忽略这点其实DNS解析可能耗时数百毫秒。解决方案是启用DNS缓存Bean public HttpClient httpClient() { return HttpClient.newBuilder() .executor(Executors.newVirtualThreadPerTaskExecutor()) .connectTimeout(Duration.ofSeconds(10)) // 关键启用DNS缓存TTL 30秒 .sslContext(SSLContext.getDefault()) .build(); }更彻底的方案是用/etc/hosts文件硬编码AI服务IP或在K8s中配置Service DNS策略。第二步检查JSON序列化性能。Jackson默认配置在大数据量时很慢。我们对比过不同配置配置项10KB JSON序列化耗时备注默认Jackson12ms启用JsonInclude(JsonInclude.Include.NON_NULL)后降至8msJackson ObjectWriter复用3ms预编译ObjectWriter避免每次反射Gson5ms无反射但功能不如Jackson丰富JSON-B (Jakarta)7ms标准化但生态弱最佳实践用ObjectWriter复用。在Configuration类中初始化Bean public ObjectMapper objectMapper() { ObjectMapper mapper new ObjectMapper(); mapper.configure(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS, false); mapper.setSerializationInclusion(JsonInclude.Include.NON_NULL); return mapper; } Bean public ObjectWriter objectWriter(ObjectMapper objectMapper) { return objectMapper.writerWithDefaultPrettyPrinter(); // 预编译Writer }然后在Service中注入ObjectWriter调用writeValueAsString()。第三步TCP连接复用。HTTP客户端默认不复用连接每次请求新建TCP连接。解决方案是启用连接池Bean public HttpClient httpClient() { // 使用Apache HttpClient连接池比JDK HttpClient更成熟 PoolingHttpClientConnectionManager connectionManager new PoolingHttpClientConnectionManager(); connectionManager.setMaxTotal(200); // 最大连接数 connectionManager.setDefaultMaxPerRoute(50); // 每路由最大连接数 CloseableHttpClient client HttpClients.custom() .setConnectionManager(connectionManager) .setKeepAliveStrategy((response, context) - 30 * 1000L) // 保持连接30秒 .build(); return HttpClientBuilder.create() .setHttpClient(client) .build(); }实测开启连接池后AI服务P95延迟从1800ms降至620ms降幅65%。4.3 “高并发下数据库写入失败怎么解决”——批量插入与分库分表策略AI服务返回结果后需要批量写入数据库。当QPS超过200时单表INSERT会成为瓶颈。我们经历过三次数据库写入失败第一次用JdbcTemplate.batchUpdate()每批10条QPS 300时MySQL CPU 100%第二次改用MyBatis-Plus的saveBatch()每批100条QPS 500时主键冲突Duplicate entry第三次引入ShardingSphere分库分表QPS 2000时稳定。根本原因是高并发下INSERT产生大量行锁且自增主键竞争激烈。解决方案分三层第一层应用层批量合并。不追求实时写入而是用ScheduledThreadPoolExecutor定时合并Component public class BatchInsertService { private final BlockingQueueAiResult resultQueue new LinkedBlockingQueue(); private final JdbcTemplate jdbcTemplate; // 每5秒批量插入一次 Scheduled(fixedDelay 5000) public void batchInsert() { ListAiResult batch new ArrayList(1000); resultQueue.drainTo(batch, 1000); // 一次最多取1000条 if (!batch.isEmpty()) { String sql INSERT INTO ai_result (task_id, result, created_time) VALUES (?, ?, ?); ListObject[] args batch.stream() .map(r - new Object[]{r.getTaskId(), r.getResult(), r.getCreatedTime()}) .collect(Collectors.toList()); jdbcTemplate.batchUpdate(sql, args); } } }第二层数据库层优化。MySQL配置关键参数# my.cnf innodb_buffer_pool_size 70% of RAM innodb_log_file_size 1G innodb_flush_log_at_trx_commit 2 # 平衡性能与安全性 innodb_autoinc_lock_mode 2 # 避免自增锁竞争第三层架构层分片。当单库写入QPS超1000必须分库。我们用ShardingSphere-JDBC按task_id哈希分片# application-sharding.yml spring: shardingsphere: props: sql-show: false rules: - !SHARDING tables: ai_result: actual-data-nodes: ds${0..1}.ai_result_${0..3} table-strategy: standard: sharding-column: task_id sharding-algorithm-name: table-inline database-strategy: standard: sharding-column: task_id sharding-algorithm-name: db-inline sharding-algorithms: db-inline: type: HASH_MOD props: sharding-count: 2 table-inline: type: HASH_MOD props: sharding-count: 4最终实现单节点QPS 2000写入延迟P9550ms。4.4 “如何监控AI服务健康度”——从日志到指标的全链路可观测性没有监控的高并发系统就像蒙眼开车。我们构建了三层监控体系第一层日志监控。用Logback配置异步Appender避免日志IO阻塞主线程!-- logback-spring.xml -- appender nameASYNC_FILE classch.qos.logback.classic.AsyncAppender appender-ref refFILE/ queueSize10000/queueSize discardingThreshold0/discardingThreshold includeCallerDatafalse/includeCallerData /appender关键日志字段必须包含task_id、ai_service_name、response_time_ms、status_code、error_type。用ELK收集后可快速统计按AI服务统计成功率SELECT service, AVG(CASE WHEN status200 THEN 1 ELSE 0 END) FROM logs GROUP BY service按时间段统计P95延迟SELECT HOUR(timestamp), PERCENTILE_CONT(0.95) WITHIN GROUP (ORDER BY response_time) FROM logs GROUP BY HOUR(timestamp)第二层指标监控。用Micrometer暴露关键指标| 指标名 | 类型 | 说明 | 报警阈值 | |
返回列表