|
@@ -16,6 +16,7 @@ import org.springframework.http.codec.ServerSentEvent;
|
|
|
import org.springframework.security.core.context.SecurityContextHolder;
|
|
import org.springframework.security.core.context.SecurityContextHolder;
|
|
|
import org.springframework.web.bind.annotation.*;
|
|
import org.springframework.web.bind.annotation.*;
|
|
|
import reactor.core.publisher.Flux;
|
|
import reactor.core.publisher.Flux;
|
|
|
|
|
+import reactor.core.publisher.SignalType;
|
|
|
|
|
|
|
|
import org.springframework.jdbc.core.JdbcTemplate;
|
|
import org.springframework.jdbc.core.JdbcTemplate;
|
|
|
import org.springframework.web.multipart.MultipartFile;
|
|
import org.springframework.web.multipart.MultipartFile;
|
|
@@ -40,13 +41,15 @@ public class ChatController {
|
|
|
private final JdbcTemplate jdbc;
|
|
private final JdbcTemplate jdbc;
|
|
|
private final QwenProperties props;
|
|
private final QwenProperties props;
|
|
|
private final HttpServletRequest request;
|
|
private final HttpServletRequest request;
|
|
|
|
|
+ private final AnalyticsService analyticsService;
|
|
|
|
|
|
|
|
public ChatController(RetrieverService rs, LLMService ls, PromptService ps,
|
|
public ChatController(RetrieverService rs, LLMService ls, PromptService ps,
|
|
|
ChatPersistenceService cps, RerankerService rrs,
|
|
ChatPersistenceService cps, RerankerService rrs,
|
|
|
QACacheService qaCache,
|
|
QACacheService qaCache,
|
|
|
BrandRecommendService brandRecommendService,
|
|
BrandRecommendService brandRecommendService,
|
|
|
JdbcTemplate jdbc, QwenProperties props,
|
|
JdbcTemplate jdbc, QwenProperties props,
|
|
|
- HttpServletRequest request) {
|
|
|
|
|
|
|
+ HttpServletRequest request,
|
|
|
|
|
+ AnalyticsService analyticsService) {
|
|
|
this.retrieverService = rs;
|
|
this.retrieverService = rs;
|
|
|
this.llmService = ls;
|
|
this.llmService = ls;
|
|
|
this.promptService = ps;
|
|
this.promptService = ps;
|
|
@@ -57,6 +60,7 @@ public class ChatController {
|
|
|
this.jdbc = jdbc;
|
|
this.jdbc = jdbc;
|
|
|
this.props = props;
|
|
this.props = props;
|
|
|
this.request = request;
|
|
this.request = request;
|
|
|
|
|
+ this.analyticsService = analyticsService;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
@PostMapping("/ask")
|
|
@PostMapping("/ask")
|
|
@@ -167,6 +171,10 @@ public class ChatController {
|
|
|
|
|
|
|
|
final String intent = retrieverService.classifyIntent(query);
|
|
final String intent = retrieverService.classifyIntent(query);
|
|
|
final long t0 = System.currentTimeMillis();
|
|
final long t0 = System.currentTimeMillis();
|
|
|
|
|
+ final String chatRequestId = request.getHeader("X-Request-Id");
|
|
|
|
|
+ final String chatIp = IpUtils.getClientIp(request);
|
|
|
|
|
+ final String chatUa = request.getHeader("User-Agent");
|
|
|
|
|
+ final java.util.Map<String, Long> timings = new java.util.concurrent.ConcurrentHashMap<>();
|
|
|
log.info("[chatStream] intent={}, query={}", intent, query.substring(0, Math.min(50, query.length())));
|
|
log.info("[chatStream] intent={}, query={}", intent, query.substring(0, Math.min(50, query.length())));
|
|
|
|
|
|
|
|
// 先发射 intent/status 事件,再用 flatMapMany 接回管道保持取消链完整。
|
|
// 先发射 intent/status 事件,再用 flatMapMany 接回管道保持取消链完整。
|
|
@@ -179,11 +187,13 @@ public class ChatController {
|
|
|
.flatMapMany(docs -> {
|
|
.flatMapMany(docs -> {
|
|
|
long t2 = System.currentTimeMillis();
|
|
long t2 = System.currentTimeMillis();
|
|
|
log.info("[chatStream] search+rerank done, docs={}, elapsed={}ms", docs.size(), t2 - t0);
|
|
log.info("[chatStream] search+rerank done, docs={}, elapsed={}ms", docs.size(), t2 - t0);
|
|
|
|
|
+ timings.put("search_rerank_ms", t2 - t0);
|
|
|
|
|
|
|
|
final List<Map<String, String>> messages = promptService.buildPrompt(query, docs, intent);
|
|
final List<Map<String, String>> messages = promptService.buildPrompt(query, docs, intent);
|
|
|
final long t3 = System.currentTimeMillis();
|
|
final long t3 = System.currentTimeMillis();
|
|
|
log.info("[chatStream] prompt built, elapsed={}ms, chars={}",
|
|
log.info("[chatStream] prompt built, elapsed={}ms, chars={}",
|
|
|
t3 - t0, messages.stream().mapToInt(m -> m.get("content").length()).sum());
|
|
t3 - t0, messages.stream().mapToInt(m -> m.get("content").length()).sum());
|
|
|
|
|
+ timings.put("prompt_ms", t3 - t0);
|
|
|
|
|
|
|
|
// 用 StringBuilder 攒完整答案(Flux 内部同步操作,无线程安全问题)
|
|
// 用 StringBuilder 攒完整答案(Flux 内部同步操作,无线程安全问题)
|
|
|
var fullAnswerBuf = new StringBuilder();
|
|
var fullAnswerBuf = new StringBuilder();
|
|
@@ -212,6 +222,8 @@ public class ChatController {
|
|
|
}
|
|
}
|
|
|
long t4 = System.currentTimeMillis();
|
|
long t4 = System.currentTimeMillis();
|
|
|
log.info("[chatStream] chatStream done, llmElapsed={}ms, total={}ms", t4 - t3, t4 - t0);
|
|
log.info("[chatStream] chatStream done, llmElapsed={}ms, total={}ms", t4 - t3, t4 - t0);
|
|
|
|
|
+ timings.put("llm_ms", t4 - t3);
|
|
|
|
|
+ timings.put("total_internal_ms", t4 - t0);
|
|
|
persistenceService.saveMessage(userKey, cid, "user", query, intent, null);
|
|
persistenceService.saveMessage(userKey, cid, "user", query, intent, null);
|
|
|
persistenceService.saveMessage(userKey, cid, "assistant", finalAnswer, intent, sources, brandRecs);
|
|
persistenceService.saveMessage(userKey, cid, "assistant", finalAnswer, intent, sources, brandRecs);
|
|
|
qaCache.put(normalized, finalAnswer, intent, sources);
|
|
qaCache.put(normalized, finalAnswer, intent, sources);
|
|
@@ -233,7 +245,8 @@ public class ChatController {
|
|
|
log.error("[chatStream] error, totalElapsed={}ms, error={}", t4 - t0, e.getMessage());
|
|
log.error("[chatStream] error, totalElapsed={}ms, error={}", t4 - t0, e.getMessage());
|
|
|
qaCache.removePending(normalized);
|
|
qaCache.removePending(normalized);
|
|
|
});
|
|
});
|
|
|
- }));
|
|
|
|
|
|
|
+ }))
|
|
|
|
|
+ .doFinally(sig -> recordChatTiming(userKey, chatRequestId, chatIp, chatUa, "stream", t0, sig, timings));
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
/** 将缓存命中结果以流式 SSE 形式返回 */
|
|
/** 将缓存命中结果以流式 SSE 形式返回 */
|
|
@@ -873,4 +886,33 @@ public class ChatController {
|
|
|
.replaceAll("\\n{3,}", "\n\n")
|
|
.replaceAll("\\n{3,}", "\n\n")
|
|
|
.trim();
|
|
.trim();
|
|
|
}
|
|
}
|
|
|
|
|
+
|
|
|
|
|
+ /** 流式问答结束时上报一次"入参→出参"耗时(含分段),用 request_id 与前端 chat_complete 关联 */
|
|
|
|
|
+ private void recordChatTiming(String userKey, String requestId, String ip, String ua,
|
|
|
|
|
+ String endpoint, long t0, SignalType sig,
|
|
|
|
|
+ java.util.Map<String, Long> timings) {
|
|
|
|
|
+ try {
|
|
|
|
|
+ long total = System.currentTimeMillis() - t0;
|
|
|
|
|
+ String status;
|
|
|
|
|
+ if (sig == null) {
|
|
|
|
|
+ status = "unknown";
|
|
|
|
|
+ } else {
|
|
|
|
|
+ status = switch (sig) {
|
|
|
|
|
+ case ON_COMPLETE -> "ok";
|
|
|
|
|
+ case ON_ERROR -> "error";
|
|
|
|
|
+ case CANCEL -> "cancelled";
|
|
|
|
|
+ default -> sig.name().toLowerCase();
|
|
|
|
|
+ };
|
|
|
|
|
+ }
|
|
|
|
|
+ java.util.Map<String, Object> ed = new java.util.LinkedHashMap<>();
|
|
|
|
|
+ ed.put("endpoint", endpoint);
|
|
|
|
|
+ ed.put("request_id", requestId != null ? requestId : "");
|
|
|
|
|
+ ed.put("total_ms", total);
|
|
|
|
|
+ ed.put("status", status);
|
|
|
|
|
+ if (timings != null && !timings.isEmpty()) ed.putAll(timings);
|
|
|
|
|
+ analyticsService.saveEvent(userKey, "chat_stream_server", ed, "", "", ip, ua, "AI药典");
|
|
|
|
|
+ } catch (Exception ignored) {
|
|
|
|
|
+ // 埋点失败不影响主流程
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
}
|
|
}
|