This commit is contained in:
luke 2025-06-23 19:26:53 +08:00
parent eb3a858681
commit 58a78d33c1
2 changed files with 27 additions and 31 deletions

View File

@ -7,6 +7,7 @@ import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Service; import org.springframework.stereotype.Service;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter; import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.io.IOException;
import java.util.Map; import java.util.Map;
@Service @Service
@ -23,32 +24,23 @@ public class LLMAppServiceImpl implements LLMAppService {
@Override @Override
public SseEmitter ask(String llmToken, String question, Map<String, Object> params) throws Exception { public SseEmitter ask(String llmToken, String question, Map<String, Object> params) throws Exception {
SseEmitter emitter = new SseEmitter(0L); // 超时 SseEmitter emitter = new SseEmitter(0L); // 超时
try {
FilteredSseOutputAdapter adapter = new FilteredSseOutputAdapter(emitter);
// 异步执行避免阻塞返回
new Thread(() -> { new Thread(() -> {
try { try {
FilteredSseOutputAdapter adapter = new FilteredSseOutputAdapter(emitter);
llmServiceFactory.current().streamAnswer(llmToken, question, params, adapter); llmServiceFactory.current().streamAnswer(llmToken, question, params, adapter);
} catch (Exception e) { } catch (Exception e) {
log.error("LLM流式调用异常", e); log.error("LLM调用异常", e);
try { try {
emitter.send(SseEmitter.event().data("{\"error\": \"LLM异常\"}")); emitter.send(SseEmitter.event().data("{\"error\": \"LLM异常\"}"));
} catch (Exception ignored) {
log.error("error", ignored);
}
emitter.completeWithError(e); emitter.completeWithError(e);
} catch (IOException ioException) {
log.warn("SSE发送错误信息失败", ioException);
}
} }
}).start(); }).start();
} catch (Exception e) {
log.error("LLM初始化异常", e);
emitter.send(SseEmitter.event().data("{\"error\": \"LLM异常\"}"));
emitter.completeWithError(e);
}
return emitter; return emitter;
} }
} }

View File

@ -15,7 +15,8 @@ public class FilteredSseOutputAdapter implements WriterAdapter {
private final SseEmitter emitter; private final SseEmitter emitter;
private static final ObjectMapper mapper = new ObjectMapper(); private static final ObjectMapper mapper = new ObjectMapper();
private Map<String, Object> lastChunk = new HashMap<>(); private final Map<String, Object> lastChunk = new HashMap<>();
private boolean closed = false;
public FilteredSseOutputAdapter(SseEmitter emitter) { public FilteredSseOutputAdapter(SseEmitter emitter) {
this.emitter = emitter; this.emitter = emitter;
@ -41,27 +42,30 @@ public class FilteredSseOutputAdapter implements WriterAdapter {
lastChunk.put("textResponse", original.get("textResponse")); lastChunk.put("textResponse", original.get("textResponse"));
} }
// 只缓存 sources不立即发送 // sources 只保存不立刻发
if (original.containsKey("sources")) { if (original.containsKey("sources")) {
lastChunk.put("sources", original.get("sources")); lastChunk.put("sources", original.get("sources"));
} }
// close 也缓存并判断是否最后一条
if (original.containsKey("close")) { if (original.containsKey("close")) {
filtered.put("close", original.get("close"));
lastChunk.put("close", original.get("close")); lastChunk.put("close", original.get("close"));
} }
// 最后一条才合并 sources 输出
if (Boolean.TRUE.equals(original.get("close"))) { if (Boolean.TRUE.equals(original.get("close"))) {
// 合并缓存中的 sources // 合并缓存数据并一次性输出
filtered.putAll(lastChunk); filtered.putAll(lastChunk);
lastChunk.clear(); // 清空缓存 lastChunk.clear();
String payload = mapper.writeValueAsString(filtered); String payload = mapper.writeValueAsString(filtered);
emitter.send(SseEmitter.event().data(payload)); emitter.send(SseEmitter.event().data(payload));
// 增加主动关闭连接逻辑
if (!closed) {
closed = true;
emitter.complete(); emitter.complete();
}
} else if (!filtered.isEmpty()) { } else if (!filtered.isEmpty()) {
// 普通中间 // 中间段
String payload = mapper.writeValueAsString(filtered); String payload = mapper.writeValueAsString(filtered);
emitter.send(SseEmitter.event().data(payload)); emitter.send(SseEmitter.event().data(payload));
} }