Files
happy-life-star/docs/superpowers/specs/2026-07-21-sse-streaming-and-history-save-fix-design.md
T

19 KiB
Raw Blame History

author, created_at, purpose
author created_at purpose
AI Assistant 2026-07-21 修复 SSE 流式传输不工作(小说生成一次性输出)和历史列表不保存最新生成小说的问题

SSE 流式传输与历史列表保存修复设计

问题概述

问题 1:历史列表不保存最新生成的小说

现象:用户完成小说生成后,历史列表页面没有显示新生成的小说。

根本原因

  • 后端 ShortNovelServiceImpl.forwardSse() 方法中,novel_done 事件触发保存数据库的条件是 originalQuery != null
  • 前端在调用 followupStream() 时没有正确传递 originalQuery 参数,或传递的是空字符串
  • 导致所有通过 followup 接口触发的小说生成(澄清回答、大纲确认/修改、重试)都不会保存到数据库

代码位置

  • 后端:server/src/main/java/com/emotion/service/impl/ShortNovelServiceImpl.java:165
  • 前端:mini-program/src/pages/main/ScriptView.vue:1768-1830

问题 2:小说生成没有真正的流式逐字输出

现象:小说生成完成后一次性显示全部内容,而不是像 ChatGPT/通义千问那样逐字显示。

根本原因

  • 后端日志显示所有 novel_delta 事件都在同一秒内到达(例如 22:29:26 的所有事件)
  • 可能的原因:
    1. nginx 缓冲了 SSE 响应(即使设置了 300s 超时,但默认启用了 proxy_buffering
    2. OkHttp 的 BufferedSource.readUtf8Line() 在某些情况下可能等待直到有足够数据
    3. SseEmitter 发送事件时没有显式 flush
  • 前端的 novel_delta 处理逻辑是正确的(直接累加 delta),但事件一次性到达导致看起来不是流式的

代码位置

  • 后端:server/src/main/java/com/emotion/service/impl/ShortNovelServiceImpl.java:144-192
  • nginx/etc/nginx/sites-enabled/lifescript.happylifeos.com.conf
  • 前端:mini-program/src/pages/main/ScriptView.vue:1715-1722

设计目标

  1. 历史列表保存:确保每次成功的小说生成都保存到数据库,并在历史列表中显示
  2. 真正流式输出:实现像 ChatGPT/通义千问那样的逐字流式输出,用户可以看到文字逐个出现
  3. 可观测性:添加详细日志便于后续问题排查和性能优化

技术方案

方案 A:彻底重构 SSE 转发逻辑(推荐)

1. 后端 OkHttp 读取优化

文件server/src/main/java/com/emotion/service/impl/ShortNovelServiceImpl.java

修改内容

// 在 forwardSse 方法中添加详细日志
log.info("[ShortNovel SSE] 开始读取上游响应: path={}, sessionId={}", path, currentUserId);

while (!source.exhausted()) {
    String line = source.readUtf8LineStrict();  // 使用严格行读取
    if (line == null) break;
    
    // 添加时间戳日志
    log.debug("[ShortNovel SSE] 读取到行: length={}, timestamp={}", 
              line.length(), System.currentTimeMillis());
    
    if (line.startsWith("data:")) {
        dataBuffer.append(line.substring(5).trim());
    } else if (line.isEmpty() && dataBuffer.length() > 0) {
        String dataStr = dataBuffer.toString();
        dataBuffer.setLength(0);
        
        if ("[DONE]".equals(dataStr)) {
            continue;
        }
        
        try {
            JSONObject event = JSON.parseObject(dataStr);
            String type = event.getString("type");
            long eventTimestamp = System.currentTimeMillis();
            
            log.info("[ShortNovel SSE] 处理事件: type={}, timestamp={}", type, eventTimestamp);
            
            // 拦截 novel_done 事件,保存到数据库
            if ("novel_done".equals(type)) {
                JSONObject payload = event.getJSONObject("payload");
                log.info("[ShortNovel SSE] novel_done 事件: originalQuery={}, payload={}", 
                        originalQuery, payload != null);
                
                if (payload != null && originalQuery != null && !originalQuery.trim().isEmpty()) {
                    String fullText = payload.getString("full_text");
                    if (fullText != null) {
                        Map<String, Object> metadata = new HashMap<>();
                        if (payload.get("title") != null) {
                            metadata.put("title", payload.get("title"));
                        }
                        
                        log.info("[ShortNovel SSE] 开始保存小说: userId={}, queryLength={}, textLength={}", 
                                currentUserId, originalQuery.length(), fullText.length());
                        
                        Map<String, String> saveResult = epicScriptDialogueServiceImpl.saveNovelResult(
                                currentUserId,
                                originalQuery,
                                fullText,
                                metadata);
                        
                        log.info("[ShortNovel SSE] 小说保存成功: scriptId={}", saveResult.get("scriptId"));
                        
                        // 注入 scriptId 到事件中
                        payload.put("scriptId", saveResult.get("scriptId"));
                        payload.put("conversationId", saveResult.get("conversationId"));
                        payload.put("currentVersionMessageId", saveResult.get("currentVersionMessageId"));
                    } else {
                        log.warn("[ShortNovel SSE] novel_done 事件缺少 full_text");
                    }
                } else {
                    log.warn("[ShortNovel SSE] novel_done 事件跳过保存: originalQuery={}, payload={}", 
                            originalQuery, payload);
                }
            }
            
            // 逐事件转发给前端,并显式 flush
            emitter.send(SseEmitter.event().name(type).data(event.toJSONString()));
            emitter.send(SseEmitter.event().comment(""));  // 触发 flush
        } catch (Exception parseEx) {
            log.warn("[ShortNovel SSE] 事件解析失败: {}", parseEx.getMessage());
        }
    }
}

log.info("[ShortNovel SSE] 完成读取上游响应");

关键点

  • 使用 readUtf8LineStrict() 替代 readUtf8Line()
    • readUtf8Line() 在遇到 EOF 时可能返回不完整的数据
    • readUtf8LineStrict() 严格遵循 SSE 规范,遇到 \n\r\n 才返回一行
    • 这确保每行数据都是完整的 SSE 事件,避免缓冲导致的数据堆积
  • 添加毫秒级时间戳日志,便于诊断事件到达的时间间隔
  • emitter.send() 后发送空注释事件,触发 SseEmitter 的 flush
    • Spring 的 SseEmitter 默认会缓冲事件,直到缓冲区满或连接关闭
    • 发送空注释 comment("") 会强制刷新缓冲区,确保事件立即发送到前端
    • 这是实现真正流式传输的关键步骤
  • 增加 originalQuery.trim().isEmpty() 检查,避免空字符串被误判为有效值

2. nginx 配置优化

文件/etc/nginx/sites-enabled/lifescript.happylifeos.com.conf

修改内容

location /api/shortNovel/ {
    proxy_pass http://127.0.0.1:19089;
    proxy_set_header Host $host;
    proxy_set_header X-Real-IP $remote_addr;
    proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
    proxy_set_header X-Forwarded-Proto $scheme;
    
    # 禁用代理缓冲,确保 SSE 事件立即转发
    proxy_buffering off;
    proxy_cache off;
    chunked_transfer_encoding on;
    proxy_set_header Connection '';
    
    # 超时配置
    proxy_connect_timeout 300s;
    proxy_send_timeout 300s;
    proxy_read_timeout 300s;
}

关键点

  • proxy_buffering off:禁用 nginx 的响应缓冲
  • proxy_cache off:禁用缓存
  • chunked_transfer_encoding on:启用分块传输编码
  • proxy_set_header Connection '':清除 Connection 头,避免干扰

3. 前端事件处理优化

文件mini-program/src/pages/main/ScriptView.vue

修改内容

// 在 handleShortNovelEvent 中添加日志
const handleShortNovelEvent = (event) => {
  const { type, session_id, payload = {} } = event
  const timestamp = Date.now()
  
  console.log('[ScriptView] 收到事件:', { type, session_id, timestamp, payload })
  
  if (session_id) novelSessionId.value = session_id
  
  switch (type) {
    case 'novel_delta': {
      const lastNovel = [...resultMessages.value].reverse().find(m => m.kind === 'novel' && m.pending)
      if (lastNovel) {
        const delta = payload.delta || ''
        console.log('[ScriptView] novel_delta:', { 
          deltaLength: delta.length, 
          currentLength: lastNovel.content.length,
          timestamp 
        })
        lastNovel.content += delta
      }
      keepResultAtBottom()
      break
    }
    case 'novel_done': {
      console.log('[ScriptView] novel_done:', { 
        scriptId: payload.scriptId, 
        hasFullText: !!payload.full_text,
        timestamp 
      })
      // 标记最后一条 novel 消息完成
      const lastNovel = [...resultMessages.value].reverse().find(m => m.kind === 'novel')
      if (lastNovel) {
        if (payload.full_text) lastNovel.content = payload.full_text
        lastNovel.pending = false
      }
      scriptId.value = payload.scriptId || ''
      conversationId.value = payload.conversationId || ''
      currentVersionMessageId.value = payload.currentVersionMessageId || ''
      generationPhase.value = 'done'
      generationStatus.value = 'idle'
      pendingNextResponse.value = false
      persistResultMessages()
      store.fetchScripts()
      break
    }
  }
}

// 在 startNovelGeneration 中确认 firstQuery 设置
const startNovelGeneration = async (source = 'home') => {
  const text = wishText.value.trim()
  if (!text) return
  
  // ... 前面的验证代码 ...
  
  firstQuery.value = text  // 确认设置首次心愿文本
  console.log('[ScriptView] 设置 firstQuery:', { text, length: text.length })
  
  // ... 后面的初始化代码 ...
}

// 在 submitClarification 中添加日志确认 originalQuery 传递
const submitClarification = (answer) => {
  if (answeringClarification.value) return
  const card = [...resultMessages.value].reverse().find(m => m.kind === 'card' && !m.submitted)
  const option = card?.card?.options?.find(opt => opt.value === answer)
  const label = option?.label || answer
  if (card) {
    card.submitted = true
    card.answer = label
  }
  pendingNextResponse.value = true
  answeringClarification.value = true
  console.log('[ScriptView] submitClarification:', { 
    sessionId: novelSessionId.value,
    originalQuery: firstQuery.value,
    originalQueryLength: firstQuery.value?.length
  })
  currentStreamTask.value = followupStream({
    sessionId: novelSessionId.value,
    action: 'answer_clarification',
    payload: { answer },
    originalQuery: firstQuery.value,
    onEvent: handleShortNovelEvent,
    onError: (errMsg) => markGenerationFailed(errMsg)
  })
  setTimeout(() => { answeringClarification.value = false }, 200)
}

// 在 confirmOutline 中添加日志和确认 originalQuery 传递
const confirmOutline = (msg) => {
  if (msg) msg.confirmed = true
  pendingNextResponse.value = true
  addResultMessage({ role: 'user', kind: 'text', content: '确认大纲,开始生成小说' })
  console.log('[ScriptView] confirmOutline:', { 
    sessionId: novelSessionId.value,
    originalQuery: firstQuery.value,
    originalQueryLength: firstQuery.value?.length
  })
  currentStreamTask.value = followupStream({
    sessionId: novelSessionId.value,
    action: 'confirm_outline',
    payload: null,
    originalQuery: firstQuery.value,
    onEvent: handleShortNovelEvent,
    onError: (errMsg) => markGenerationFailed(errMsg)
  })
}

// 在 modifyOutline 中添加日志和确认 originalQuery 传递
const modifyOutline = (msg, feedback) => {
  if (!feedback.trim()) {
    uni.showToast({ title: '请先填写修改意见', icon: 'none' })
    return
  }
  if (msg) msg.confirmed = true
  pendingNextResponse.value = true
  addResultMessage({ role: 'user', kind: 'text', content: feedback })
  console.log('[ScriptView] modifyOutline:', { 
    sessionId: novelSessionId.value,
    originalQuery: firstQuery.value,
    originalQueryLength: firstQuery.value?.length
  })
  currentStreamTask.value = followupStream({
    sessionId: novelSessionId.value,
    action: 'modify_outline',
    payload: { feedback },
    originalQuery: firstQuery.value,
    onEvent: handleShortNovelEvent,
    onError: (errMsg) => markGenerationFailed(errMsg)
  })
}

// 在 resumeSession 中添加日志和确认 originalQuery 传递
const resumeSession = () => {
  if (!resumeableSessionId.value || generating.value) return
  novelSessionId.value = resumeableSessionId.value
  resumeableSessionId.value = ''
  generating.value = true
  generationPhase.value = 'generating'
  generationStatus.value = 'waiting'
  generationError.value = ''
  addResultMessage({ role: 'user', kind: 'text', content: '继续之前的创作' })
  streamWriter.reset()
  startGenerationFeedback()
  keepResultAtBottom()
  console.log('[ScriptView] resumeSession:', { 
    sessionId: novelSessionId.value,
    originalQuery: firstQuery.value,
    originalQueryLength: firstQuery.value?.length
  })
  
  try {
    currentStreamTask.value = followupStream({
      sessionId: novelSessionId.value,
      action: 'retry',
      payload: null,
      originalQuery: firstQuery.value,
      onEvent: handleShortNovelEvent,
      onError: (errMsg) => {
        markGenerationFailed(errMsg)
        analytics.track('script_generate_fail', {
          source: 'resume',
          error: errMsg
        }, { eventType: 'script', pagePath })
      }
    })
  } catch (error) {
    markGenerationFailed(error.message || '恢复失败')
  } finally {
    generating.value = false
  }
}

关键点

  • 添加详细的 console.log,记录每个事件的接收时间
  • 确认 firstQuery.value 在首次请求时正确设置
  • 在所有 followupStream 调用中传递 originalQuery

方案 B:诊断优先(备选)

如果方案 A 不能解决问题,可以采用方案 B:

  1. 添加更详细的诊断日志

    • 后端:记录每个事件的接收时间、发送时间、事件类型、数据长度
    • 前端:记录每个事件的接收时间、事件类型、数据长度
    • nginx:启用 access log,记录请求和响应的时间戳
  2. 分析日志确定瓶颈

    • 后端接收事件的时间间隔
    • 后端发送事件的时间间隔
    • 前端接收事件的时间间隔
  3. 根据诊断结果修复

    • 如果是上游服务问题,考虑联系上游服务提供者
    • 如果是 nginx 问题,调整 nginx 配置
    • 如果是 OkHttp 问题,考虑更换 HTTP 客户端

方案 C:绕过代理层(不推荐)

如果代理层确实无法解决缓冲问题,可以考虑:

  1. 前端直接连接上游 SSE 服务(绕过 nginx 和后端)
  2. 但这会破坏现有架构,降低安全性,不利于统一监控

不推荐此方案,除非方案 A 和 B 都失败。

实施步骤

第一步:后端修改

  1. 修改 ShortNovelServiceImpl.java

    • 添加详细日志
    • 使用 readUtf8LineStrict()
    • emitter.send() 后添加 flush 触发
    • 增加 originalQuery 的空字符串检查
  2. 编译后端:

    cd server
    mvn clean install -DskipTests
    

第二步:nginx 配置

  1. 修改 /etc/nginx/sites-enabled/lifescript.happylifeos.com.conf

    • 添加 proxy_buffering off;
    • 添加 proxy_cache off;
    • 添加 chunked_transfer_encoding on;
    • 添加 proxy_set_header Connection '';
  2. 重新加载 nginx

    sudo nginx -t
    sudo systemctl reload nginx
    

第三步:前端修改

  1. 修改 ScriptView.vue

    • 添加 console.log 日志
    • 确认 firstQuery.value 正确设置
    • 确认所有 followupStream 调用传递 originalQuery
  2. 构建小程序:

    cd mini-program
    npm run build:mp-weixin
    

第四步:部署

  1. 上传后端 JAR 到服务器:

    scp server/target/server-1.0.0.jar root@101.200.208.45:/data/programs/emotion-museum/
    
  2. 重启后端服务:

    ssh root@101.200.208.45 "cd /data/programs/emotion-museum && ./deploy-server.sh test"
    

第五步:验证

  1. 后端验证

    • 查看日志,确认事件逐条到达(时间间隔 > 100ms)
    • 确认 novel_done 事件触发保存,日志显示 小说保存成功
  2. 前端验证

    • 浏览器 Console 查看日志,确认事件逐条到达
    • 观察小说生成是否逐字显示
  3. 数据库验证

    • 查询 t_epic_script 表,确认新记录存在
    • 检查历史列表页面,确认显示新生成的小说

测试用例

测试用例 1:完整生成流程

  1. 输入心愿文本
  2. 回答澄清问题
  3. 确认大纲
  4. 等待小说生成完成
  5. 验证:
    • 前端逐字显示小说内容
    • 后端日志显示事件逐条到达
    • 数据库中出现新记录
    • 历史列表显示新生成的小说

测试用例 2:修改大纲后重新生成

  1. 输入心愿文本
  2. 回答澄清问题
  3. 修改大纲
  4. 等待小说生成完成
  5. 验证:同上

测试用例 3:继续之前的创作

  1. 输入心愿文本
  2. 回答澄清问题
  3. 中断流程(关闭页面)
  4. 重新进入页面,点击"继续创作"
  5. 等待小说生成完成
  6. 验证:同上

风险评估

风险 1:上游服务不支持真正的流式传输

可能性:中等

影响:即使后端和 nginx 都配置正确,如果上游服务一次性发送所有数据,流式输出仍然不会工作

缓解措施

  • 添加详细的日志,记录每个事件的接收时间
  • 如果日志显示所有事件在同一毫秒到达,说明是上游服务问题
  • 联系上游服务提供者,确认是否支持真正的流式传输

风险 2nginx 配置不生效

可能性:低

影响:nginx 仍然缓冲响应,导致流式输出不工作

缓解措施

  • 使用 nginx -t 验证配置语法
  • 使用 curl -N 测试 SSE 接口,观察是否逐条接收事件
  • 检查 nginx 的 error log,确认没有配置错误

风险 3SseEmitter flush 不工作

可能性:低

影响:即使调用了 emitter.send(),事件仍然被缓冲

缓解措施

  • 使用 emitter.send(SseEmitter.event().comment("")) 触发 flush
  • 如果仍然不工作,考虑更换为 ResponseBodyEmitterStreamingResponseBody

成功标准

  1. 流式输出:小说生成过程中,前端逐字显示内容,用户可以观察到文字逐个出现
  2. 历史列表保存:每次成功的小说生成都保存到数据库,并在历史列表中显示
  3. 日志可观测:后端和前端日志清晰记录每个事件的处理过程,便于问题排查

后续优化

  1. 性能监控:添加 Prometheus 指标,监控 SSE 连接的延迟和吞吐量
  2. 错误重试:如果 SSE 连接中断,支持自动重连
  3. 离线缓存:对于生成失败的小说,支持离线缓存,网络恢复后重试