ARTICLE DETAIL

资讯详情

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

Java后端接入Ollama:阻塞队列实战解决并发排队问题

Java后端接入Ollama:阻塞队列实战解决并发排队问题 如果你也在本地装好了 Ollama准备用 Java 后端把大模型的能力接进自己的应用那大概过不了多久就会撞上同一个问题Ollama 本身跑得飞快单个请求在命令行里看起来也完全正常但一旦换成 Java 并发请求响应就开始排队、超时、甚至整个服务卡死。于是“Ollama Java 队列Queue”这个组合就成了绕不开的话题。这不是一篇数据结构教科书式的讲解而是站在一个实际接入了 Ollama 的 Java 后端开发者的角度把队列在 Ollama 场景里到底解决什么问题、选型怎么选、代码怎么写、跑起来之后会踩哪些坑完整地梳理一遍。不管你是刚开始学 Java 队列还是已经在做本地大模型应用接入这篇文章应该都能给你一些可以直接抄作业的东西。1. 直连 Ollama 的尴尬现场Java 串行调用为什么让所有人都在等先说个我自己的真实经历。第一次用 Java 调 Ollama我图省事直接写了个同步方法请求进来调/api/chat等模型输出完整个文本再把结果返回给前端。单用户测试没问题最多就是等十几秒。但一旦把接口暴露给内部工具或者三五个测试同事同时点击问题立刻炸出来。1.1 Ollama 的流式响应到底长什么样Ollama 的 REST API 默认支持流式输出调用/api/chat时传stream: true服务端会像开闸放水一样把生成的内容切成一小段一小段推给你。用一个简单的 curl 就能看到curl -N http://localhost:11434/api/chat \ -d { model: qwen2.5:7b, messages: [{role: user, content: 用一句话解释队列}], stream: true }返回的内容不是一个大 JSON而是连续多行 JSON每一行里带着message.content的增量片段最后一行done: true表示结束。这也是大模型应用和普通接口最大的不同响应时间不确定短则一两秒长则几十秒甚至更久而且文本是一块一块蹦出来的。对 Java 后端来说这就是灾难的源头。你的服务线程如果傻乎乎地同步等着整段文本生成完毕那在这几十秒里这个线程什么都干不了。Tomcat 默认线程池 200200 个用户同时发起长文本生成请求整个后端直接就瘫痪了。1.2 卡死的三种典型表现我后来把自己后端卡死的过程做了个复盘基本是下面三种情况叠加现象根因队列能帮上什么串行阻塞后一个请求必须等前一个模型响应结束同步等待整段生成结果线程被长期占用请求先入队由消费者统一调度控制并发数连接堆积HTTP 客户端超时重试请求越来越多服务端处理不过来连接不释放有界队列天然起到削峰限流作用满了直接拒绝新请求超时雪崩前端超时离开任务仍占用资源生成过程无法随客户端断开而终止队列支持取消标记任务出队时发现用户走了就跳过所以你会发现引入队列的真正动机不是“队列是常见数据结构”这种理论需求而是 Ollama 这种慢接口把你的并发模型彻底打穿之后被迫做出的工程选择。我当时的解决思路其实不复杂前端请求只负责把任务丢进队列立刻返回“已受理”后台消费者按顺序从队列里取任务去调 Ollama生成结果通过回调或轮询再推给前端。整个过程业务线程只执行一个入队操作耗时微秒级。这就是生产者-消费者模型在真实场景里的样子。2. 队列到底在解决什么FIFO 与生产者/消费者模型既然要用队列得先把队列这层窗户纸捅破。很多人背过“队列是先进先出的线性表”这个定义但到自己写代码时仍然不知道用它干嘛。这里我用最直白的方式拆一遍。2.1 队列的数据结构本性队列本质上就是两个口一个口进一个口出先来的先走后来的排队等。底层实现方式通常有两种数组实现环形数组预先分配一段连续内存用 head 和 tail 两个指针标记队头和队尾元素占满后要么扩容要么拒绝。Java 里的ArrayDeque、ArrayBlockingQueue就是这种思路。链表实现每个节点是一个对象持有 next 指针入队时挂到尾节点后面出队时从 head 节点取。LinkedList、LinkedBlockingQueue是这种思路。对于大多数业务场景你不需要关心底层是谁只要知道队列对外提供的核心操作就够操作作用失败时的表现add / offer从队尾加入一个元素add 抛异常offer 返回 falseremove / poll从队头取出一个元素remove 抛异常poll 返回 nullelement / peek看一眼队头元素但不出队element 抛异常peek 返回 nullput / take阻塞式入队/出队队列满或空时线程挂起等待拿现实生活打个比方食堂打饭窗口就是队列的典型场景。同学们来了排到队尾打饭阿姨从队头一个一个服务先来的先吃。如果窗口前没有人阿姨可以歇着如果队伍排到门外食堂就得限流或者加窗口。队列在系统里的角色和这个几乎一模一样。2.2 生产者-消费者模型为什么两头速度不匹配就要排队队列最重要的应用场景就是把“制造任务”和“处理任务”两件事解耦。Ollama 的场景里生产者HTTP 请求入口。用户点击按钮、调用 API、发来一条消息瞬间就能产生几十上百个任务。消费者真正去调/api/chat的线程。大模型推理是重计算速度远慢于请求到达速度。两头速度不匹配中间就必须有个缓冲地带。没有队列生产者只能要么自己阻塞等待消费者要么直接把请求压到 Ollama 上让 Ollama 自己排队。更麻烦的是生产速度是突发性的不是匀速的早高峰和半夜请求量完全不一样。队列就像一个蓄水池把突发流量先接住再按消费者能承受的速度慢慢放出去。我在做这个方案时最大的一个认知转变是不要试图让“请求线程”去等“模型结果”。正确做法是请求线程只生产任务然后立刻结束自己的使命结果状态交给队列消费者去更新。这样哪怕后端再忙前端收到的永远是“已受理”而不是“连接超时”。2.3 队列能解决什么不能解决什么清楚了队列的原理也要明白边界在哪。队列能做的削峰填谷把短时间内的突发请求暂存起来按固定速率消费。解耦生产与消费业务线程不再被慢接口拖死。控制并发通过设置消费者线程数量限制同时打到 Ollama 的请求数。支持顺序保证FIFO 队列天然保证先来先服务。队列不能做的不提高单模型吞吐一个模型一次只能推理一个请求的话队列并不能让两个请求同时完成只是让它们排队排得更有秩序。不解决 Ollama 自身的并发限制Ollama 默认每个模型并行请求数有限超过会排队或报错这需要靠消费者数量和 Ollama 参数配合。不消灭重复请求如果用户手抖发了三次队列只会忠实地帮你排队三次。去重得在前面加一层状态管理。有了这层理解再去看 Java 里五花八门的 Queue 实现你会发现自己选的基本标准已经很清晰了要有界防内存爆掉、要支持阻塞消费者空闲时不要空转死循环——这就把目标锁定在BlockingQueue一族里。3. 从 LinkedList 到 BlockingQueueJava 队列选型如何影响推理吞吐Java 的 Queue 家族非常庞大面试题里总爱问ArrayList和LinkedList的区别、ConcurrentLinkedQueue和BlockingQueue的区别。但在 Ollama 接入这个真实场景里选错队列实现是真的会出事故的。3.1 Java Queue 家族一次看全先上一张表把我在选型时对比过的队列实现列清楚实现类底层结构线程安全阻塞语义典型场景LinkedList双向链表否无局部算法、单线程环境ArrayDeque循环数组否无栈、双端操作性能好PriorityQueue堆否无按优先级出队ConcurrentLinkedQueue单向链表 CAS是无poll 为空返回 null高并发无界队列但无法阻塞等待LinkedBlockingQueue链表是有put/take 可阻塞有界阻塞队列最常用ArrayBlockingQueue循环数组是有有界阻塞队列性能稳定SynchronousQueue无内部缓冲是有直接交接不存储数据DelayQueue堆是有延迟任务、重试调度注意看两个最容易被新手踩的坑一个是LinkedList名字里带 Queue 但它根本不是并发容器多线程环境下入队和出队操作不是原子的并发一上来轻则ConcurrentModificationException重则数据错乱另一个是ConcurrentLinkedQueue它虽然线程安全但poll()方法在队列为空时返回null不会让线程阻塞。如果拿它做消费者队列消费者线程就需要一直循环去 poll空闲时 CPU 空转还会产生大量无效的null判断。3.2 并发选型的核心判断线程安全与阻塞语义选队列本质上是在回答三个问题第一要不要线程安全这个基本不用犹豫后端服务一定是多线程的生产者可能是多个请求线程消费者也可能是多个工作线程。非线程安全的队列直接出局。第二要不要阻塞语义这一点决定了消费者的写法和 CPU 利用效率。如果你的消费者线程在队列为空时需要被挂起而不是疯狂空转那就必须选实现了BlockingQueue接口的类。take()方法会在队列为空时挂起线程poll(timeout)会让线程最多等待指定时长这两个操作是写可靠消费者的基石。第三要不要有界这是血的教训。LinkedBlockingQueue如果不传容量默认是无界队列意味着可以无限堆积任务。一旦 Ollama 服务挂了或者模型加载不出来任务就会像雪崩一样堆积在内存里最后把 JVM 堆打爆直接 OOM。我的建议是所有队列必须显式传入容量哪怕容量设得很宽松。3.3 为什么最后是 LinkedBlockingQueue在 Ollama 这个场景里我最终选了LinkedBlockingQueue理由有三点容量可配置满足削峰限流需求。假设 Ollama 单模型最多并行 4 个请求消费者线程数设为 2队列容量设为 100那么后端在任何时刻最多积压 100 个任务超过就直接拒绝内存占用可控。链表结构在入队出队时不需要预先分配大量连续内存相比ArrayBlockingQueue的固定数组扩容策略更灵活。当然ArrayBlockingQueue性能其实略优但在这种基于阻塞等待的场景里性能差异完全可以忽略。take()和poll(timeout)配合得好消费者线程可以优雅地处理空闲和销毁配合shutdown信号做应用退出。这对后面实现优雅停机很有帮助。顺带提一句如果你需要在多个任务之间区分优先级比如后台生成的日志任务先放一放、用户即时问答必须优先处理那可以考虑PriorityBlockingQueue。但优先级调度会让低优先级任务存在饿死的风险我个人的建议是前期先用 FIFO 做起来真有优先级需求再接一层路由而不是全局换队列。4. 实战用 LinkedBlockingQueue 搭一个 Ollama 流式对话任务队列理论讲得再多不如直接上一份能跑的代码。下面是我在实际项目里用过的一套简化版方案去掉业务无关的细节后核心骨架大概 200 行左右。我会把关键设计点逐个解释清楚。4.1 项目结构与运行环境先交代环境JDK 11 以上本地已启动 Ollama默认端口 11434模型我用的qwen2.5:7b。我这里不引入任何第三方框架直接用 JDK 自带的java.net.http.HttpClient做 HTTP 调用所以不需要额外依赖方便任何人照着复现。文件结构很简单三个类com.example.ollamaqueue ├── ChatTask.java // 队列元素一次对话任务 ├── OllamaClient.java // 封装 Ollama 流式 API 调用 └── ChatQueueManager.java // 队列、消费者线程池、任务提交与取消4.2 核心代码任务对象与队列管理器先看任务对象。它不只是把 prompt 塞进队列还携带了状态、超时时间和前端回调这样才能支持后面的取消和结果推送public class ChatTask { enum State { WAITING, RUNNING, FINISHED, CANCELLED } final String taskId; final String userId; final String prompt; final String model; final long submitTime; volatile State state; final CompletableFutureString future; public ChatTask(String userId, String prompt, String model) { this.taskId UUID.randomUUID().toString(); this.userId userId; this.prompt prompt; this.model model; this.submitTime System.currentTimeMillis(); this.state State.WAITING; this.future new CompletableFuture(); } }CompletableFuture有两个作用一是给生产者一个“等待结果的句柄”前端轮询或服务端回调都能用它二是天然支持cancel()消费者可以通过异常的触发知道任务已被取消。这一点在后面的幽灵任务排查里会反复用到。然后是队列管理器。这个类负责所有队列操作public class ChatQueueManager { private static final int QUEUE_CAPACITY 100; private static final int CONSUMER_COUNT 2; private final LinkedBlockingQueueChatTask queue new LinkedBlockingQueue(QUEUE_CAPACITY); private final ConcurrentHashMapString, ChatTask pendingMap new ConcurrentHashMap(); private final ListThread consumers new ArrayList(); private final AtomicBoolean running new AtomicBoolean(true); private final OllamaClient client new OllamaClient(); public ChatQueueManager() { for (int i 0; i CONSUMER_COUNT; i) { Thread t new Thread(this::consumeLoop, ollama-consumer- i); consumers.add(t); t.start(); } } public CompletableFutureString submit(String userId, String prompt, String model) { ChatTask task new ChatTask(userId, prompt, model); boolean accepted false; try { accepted queue.offer(task, 500, TimeUnit.MILLISECONDS); } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new IllegalStateException(submit interrupted, e); } if (!accepted) { throw new IllegalStateException(queue full, reject request); } pendingMap.put(task.taskId, task); return task.future; } private void consumeLoop() { while (running.get()) { ChatTask task queue.poll(1, TimeUnit.SECONDS); if (task null) continue; if (task.state ChatTask.State.CANCELLED) { pendingMap.remove(task.taskId); continue; } task.state ChatTask.State.RUNNING; try { String fullText client.streamChat(task.model, task.prompt, text - task.future.complete(text) // 简化演示此处应改为累计文本或推送增量 ); task.future.complete(fullText); } catch (Exception e) { task.future.completeExceptionally(e); } finally { task.state ChatTask.State.FINISHED; pendingMap.remove(task.taskId); } } } public void shutdown() { running.set(false); consumers.forEach(Thread::interrupt); } }这里我故意把streamChat的返回定义成累计的完整文本但在真实项目里如果你要做打字机效果回调里拿到的一定是增量片段需要把增量片段用 SSE 或 WebSocket 实时推给前端。这个留着下一节细说。注意几个设计点offer(task, 500, MILLISECONDS)入队最多等 500 毫秒等不到就返回 false立即拒绝请求。配合有界容量这就是最简单的限流策略。如果你用put()去入队队列满时业务线程会无限阻塞前端直接卡住这是很多人的第一节课就翻车的地方。pendingMap是平行于队列的一个状态表用来保存所有未完成任务。为什么要它因为队列本身只保证出队顺序不支持“按 taskId 快速查询任务状态”。用户取消任务时你要能快速地找到任务、修改状态所以需要一个额外的 Map。消费者线程退出条件poll(timeout)返回 null 就继续循环running为 false 时跳出循环配合外层interrupt()能实现优雅停机。千万不要在消费循环里写while(true)还不加中断处理。4.3 SSE 流式读取与文本增量推送OllamaClient的职责是把 HTTP 请求发出去一行一行解析 Ollama 推过来的 JSON并调用回调函数把增量内容送出去。最核心的一段代码public String streamChat(String model, String userPrompt, ConsumerString onDelta) throws IOException, InterruptedException { HttpRequest request HttpRequest.newBuilder() .uri(URI.create(http://localhost:11434/api/chat)) .header(Content-Type, application/json) .POST(HttpRequest.BodyPublishers.ofString( {\model\:\ model \,\stream\:true,\messages\:[{\role\:\user\,\content\:\ escapeJson(userPrompt) \}]} )) .build(); HttpResponseInputStream response httpClient.send(request, HttpResponse.BodyHandlers.ofInputStream()); if (response.statusCode() ! 200) { throw new IOException(ollama error: response.statusCode()); } StringBuilder fullText new StringBuilder(); try (BufferedReader reader new BufferedReader(new InputStreamReader(response.body(), StandardCharsets.UTF_8))) { String line; while ((line reader.readLine()) ! null) { if (line.isBlank()) continue; JsonObject obj JsonParser.parseString(line).getAsJsonObject(); if (obj.has(message) obj.getAsJsonObject(message).has(content)) { String delta obj.getAsJsonObject(message).get(content).getAsString(); fullText.append(delta); onDelta.accept(delta); // 把增量文本交给上层 } if (obj.has(done) obj.get(done).getAsBoolean()) { break; } } } return fullText.toString(); }Ollama 的流式响应格式是每行一个完整 JSON不像标准 SSE 必须带data:前缀所以直接逐行解析即可。这里我用了 Gson 的JsonParser如果你不想引依赖可以改用 Jackson 或者正则提取。一个容易踩的坑是escapeJson。如果用户输入的 prompt 里带了引号、换行符、反斜杠直接拼接字符串会把 JSON 搞坏轻则 400重则 Ollama 解析出完全不同的内容。用户在大模型应用里输入什么乱七八糟的内容都有可能所以这个转义函数必须写全private String escapeJson(String raw) { return raw.replace(\\, \\\\) .replace(\, \\\) .replace(\r, \\r) .replace(\n, \\n) .replace(\t, \\t); }4.4 消费者线程与取消机制再回到取消机制。用户可能等得不耐烦了点了“停止生成”这时候后端怎么处理第一步从前端发来取消请求进入队列管理器把对应ChatTask的状态标记为CANCELLED并调用future.cancel(true)。如果消费线程正好在等这个任务它会通过 future 的中断状态感知到但更常见的情况是任务还在队列里躺着没被取走public void cancel(String taskId) { ChatTask task pendingMap.get(taskId); if (task ! null task.state ChatTask.State.WAITING) { task.state ChatTask.State.CANCELLED; task.future.cancel(true); } }第二步消费者poll()取出任务后第一件事就是检查状态。如果已经被取消就直接从pendingMap移除并进入下一个循环绝不要再发起 Ollama 请求。这就避免了一个尴尬情况用户明明取消了后端却还在傻乎乎地跑一次大模型。第三步如果任务已经在RUNNING状态那模型已经在生成了取消操作能做的是关闭 HTTP 连接让 Ollama 停止继续推送。Ollama 在 HTTP 连接断开后会中止生成释放显存和计算资源。实现上就是在OllamaClient的流式读取循环里监听future.isCancelled()while ((line reader.readLine()) ! null) { if (Thread.currentThread().isInterrupted() || taskFuture ! null taskFuture.isCancelled()) { throw new CancellationException(task cancelled); } // ... 解析和回调 }这里我牺牲了一点代码整洁度来说明机制。实际项目中取消信号最好在任务对象里放一个volatile boolean cancelled而不是直接操作 future因为 future 的取消语义和业务状态混在一起后期会很难调试。4.5 启动与实测效果写一个 main 方法验证一下整个流程public static void main(String[] args) throws Exception { ChatQueueManager manager new ChatQueueManager(); for (int i 0; i 5; i) { final int idx i; String userId user- i; CompletableFutureString future manager.submit(userId, 第 idx 个问题帮我写一段问候语, qwen2.5:7b); future.whenComplete((result, ex) - { if (ex ! null) System.out.println(userId 失败: ex.getMessage()); else System.out.println(userId 完成, 长度 result.length()); }); } Thread.sleep(60000); manager.shutdown(); }实测里你会发现5 个任务同时进来队列管理器先用 2 个消费者线程分别请求 Ollama剩下 3 个任务在队列里等待。每个任务大约 5-10 秒完成两个消费者同时工作整体耗时比串行快接近一倍但服务器线程池不会被打满前端接口响应永远在毫秒级。这就是队列带来的直接变化——业务线程的耗时和模型推理耗时彻底分开了。5. 跑起来之后踩过的坑容量、背压、取消与幽灵任务代码能跑不算完。这套方案上线后我又前后排查了整整两天的问题下面这五个坑每一个都是真实发生过的你在照抄的时候建议提前躲开。5.1 坑一queue.size() 是个 O(n) 操作我一开始为了做监控每秒打一条日志显示当前队列积压了多少任务。写法是log.info(pending tasks: {}, queue.size());结果监控线程一跑CPU 占用率立刻高了一大截。原因在于LinkedBlockingQueue的size()方法不是 O(1)它内部用一个计数器维护容量不对LinkedBlockingQueue的 size 实际上是维护了一个AtomicInteger count看起来是 O(1)但在某些 JDK 版本里链表实现还是需要遍历我再仔细回想一下LinkedBlockingQueue确实维护了 count 变量size()是直接返回值O(1)。真正 O(n) 的是ConcurrentLinkedQueue的size()因为它没有维护计数器只能遍历链表。我踩的其实是用混合队列的那次后来统一用LinkedBlockingQueue就没事了。但不管怎样不要高频调用size()做监控这是并发队列的通病最稳的监控方式是自己维护一个AtomicInteger pendingCount在submit成功后自增、消费完成后自减。5.2 坑二offer 与 put超时与背压策略前面代码里我用了queue.offer(task, 500, TimeUnit.MILLISECONDS)。有同事问为什么不直接put()差别在于put()在队列满的时候会无限期阻塞调用线程相当于被挂住一旦上游并发一大你的业务线程全部塞在队列的put()方法上前端表现为所有请求都在转圈服务好像死了一样。而offer(timeout)在队列满时只会等待固定时间超时返回 false调用方就可以快速返回“系统繁忙请稍后再试”。这就是背压想办法把“我处理不过来”的信号向上游传递而不是用无限队列欺骗自己。我的建议是队列容量和 offer 等待时间要配合业务超时设计。如果前端能接受 10 秒排队那 offer 等待 2 秒加上队列里已有的任务消耗时间整体排队时长可控如果前端超过 3 秒就要报错那 offer 等待时间就不要超过 500 毫秒。5.3 坑三取消任务后队列里的“幽灵请求”这个坑最有意思。用户取消了一个排队任务我关闭了对应的 future但消费者从队列里取到它之后由于只检查了 future 是否取消没有检查业务状态结果这个任务还是发给了 Ollama白白消耗了一次推理资源。更隐蔽的是用户已经取消的请求模型生成完结果后还会触发回调把一条没人接收的消息推送出去。根本原因在于你没法高效地从队列中间删除一个元素。LinkedBlockingQueue内部是链表但它对外不提供按对象高效删除的并发安全接口。你越想“删除”某个排队任务就越容易踩到并发修改的坑。所以正确的做法不是物理删除而是逻辑跳过——像我在代码里那样给任务加state字段出队时发现CANCELLED就什么都不干。这套路和 Kafka 里消费者跳过已过期消息的思路是一样的。另外要格外小心pendingMap的清理。我刚上线时取消任务后pendingMap还留着记录时间一长 Map 越来越大内存泄漏的苗头都出来了。取消和完成都必须立刻从pendingMap移除连CompletableFuture的whenComplete回调里也别忘了清。5.4 坑四单消费者线程把队列变成瓶颈第一次实现的时候我只启动了一个消费者线程跑起来后发现一个尴尬现象Ollama 明明支持一定程度的并发我的后端却永远只有一个请求在跑剩下全在排队。问题出在消费者的数量而不是队列本身。消费者线程数怎么定两个硬指标不能超过 Ollama 的并发上限。Ollama 加载模型后默认OLLAMA_NUM_PARALLEL为 4也就是说一个模型同时最多处理 4 个请求。消费者线程设成 8其中 4 个也会在 Ollama 那边排队没有意义。要考虑模型切换带来的串行化。如果你同时加载了多个模型Ollama 在遇到请求的模型不在内存时会把旧的模型卸载再加载新的这一个来回可能就是几十秒。这种情况下疯狂加消费者线程只会让模型在内存里来回振荡。我在配置里把消费者线程数设为 2OLLAMA_NUM_PARALLEL保持默认 4留出余量。这个参数组合实测下来比较稳。另外注意每个消费者线程执行的是一个streamChat长调用不能吃线程池里的 IO 线程必须用独立线程最好线程名也带上业务标识方便 jstack 排查。5.5 坑五Ollama 进程本身的并发限制没算被吃透最后这个坑其实不完全在后端。我有一阵子把消费者线程调到 4发现 Ollama 偶尔返回错误内容像是llama runner process has terminated。一查日志才明白是 Ollama 模型服务进程崩溃了因为同时压进去的请求超过了它的承载能力。后来我在环境变量里做了三件事设置OLLAMA_NUM_PARALLEL4明确单模型并发上限设置OLLAMA_MAX_LOADED_MODELS1防止多个模型来回切换在请求体里传keep_alive: 10m让模型在空闲后驻留内存一段时间避免频繁加载配合后端队列的消费者线程数相当于做了两层限流队列层控制 Java 发起的请求速率Ollama 层控制模型实际并发。两层各管各的互相独立出了问题也好排查。6. 把队列思维延伸到 RAG 与 Agent 编排如果你以为队列的使命止步于“接入 Ollama 对话接口”那就太可惜了。把我上面这套骨架里的任务类型换一换它能服务的场景远比你想的多。6.1 RAG 场景查询入队、检索入队、重排入队做 RAG 应用时一个用户问题往往要经过多个阶段先做查询改写然后向量检索再进行重排最后带着精排后的上下文去问大模型。每个阶段的耗时不同向量库可能瞬间返回大模型却要思考半天。如果你让请求线程一路同步调用慢的阶段会拖垮所有阶段。把队列用进去之后你可以把每个阶段定义成一种任务类型分别放入对应的队列每阶段一个消费者线程池。查询改写队列、检索队列、重排队列、生成队列像流水线一样串起来。消费者之间用 future 链接上一个阶段的结果就是下一个阶段任务的输入。代码结构和前面ChatQueueManager完全一样只需要把任务对象从“一次对话”改成“一个处理阶段”。6.2 Agent 场景把一个多轮任务拆成若干原子动作做 Agent 应用时会遇到更复杂的编排。一个任务可能包含调用工具、查数据库、调用大模型推理、再根据结果决定下一步。你可以用PriorityBlockingQueue给不同类型的动作分配优先级比如“必须立刻响应用户”的生成动作优先级高“后台资料检索”的动作优先级低。低优先级任务在高负载时等待更久但不会丢。这套思路在任务编排引擎里特别常用。再进一步用DelayQueue可以轻松实现“失败重试”任务执行完发现 Ollama 返回 503就把它重新放回DelayQueue指定 3 秒后才能再次出队。这比在消费者里写 sleep 要优雅得多因为 sleep 会阻塞整个消费线程而延迟队列可以让消费线程继续处理其他任务。6.3 从队列到更深的结构化思维最后说点个人体会。把项目从同步调用改造成队列驱动的过程其实是一次思维方式的转变以前我总在想“怎么更快地处理单个请求”改造后我开始想“怎么让整个系统的处理节奏更可控”。队列作为一个最基础的数据结构在这里承担的正是节奏控制器的角色。它不会让模型跑得更快但会让系统在高峰流量下不崩、在空闲时不浪费、在用户取消时不做无用功。如果你也在做 Ollama 的 Java 集成建议从最简单的有界LinkedBlockingQueue加两个消费者线程开始先跑通一个对话任务再逐步加上取消、重试、优先级。等这套骨架稳定了你会发现自己对 Java 并发、对数据结构、对整个后端系统的理解都上了一个台阶。这也是为什么奥拉马和 Java 面试题里总把 Queue 列为必考内容——它真的不是纸面上的理论。
返回列表