
在前四天中,我們完成了 Generative UI 的四塊核心拼圖:
root + elements Map)。今天我們要將這些靜態零件串成一條動態的即時串流管線(Streaming Pipeline)。
在現代 AI 產品體驗中,使用者無法忍受點擊送出後面對 8 秒鐘的空白等待。我們需要讓使用者親眼看見:
status: analyzing)plan: [...])step: running -> done)chunk 漸進渲染)要達到這個效果,我們必須攻克兩大工程難題:為什麼不能用瀏覽器原生的 EventSource? 以及 未閉合的「半截 JSON(Partial JSON)」該如何安全渲染而不崩潰?
EventSource 的先天缺陷:
new EventSource(url) 僅支援 HTTP GET 請求,無法在 Request Body 中傳遞複雜的自然語言查詢、篩選條件或使用者 Context。fetch() 搭配 ReadableStream 手動解析 SSE 協定。{"title": "近一年,括號與引號尚未閉合。JSON.parse() 會立即拋出 SyntaxError;若暴力略過,前端則會一直等待到最後一秒,失去串流生長的動態體驗。
fetch + ReadableStream 建立 SSE 連線。status: 意圖分析(如 ANALYZING)plan: A* 規劃出的思維鏈步驟step: 單一步驟完成狀態與耗時chunk: Spec JSON 增量片段complete: 收尾指標與 Token 成本為了同時滿足「可觀測性進度」與「UI 漸進渲染」,一條標準連線包含以下 7 類 Typed Events:
| 事件名稱 (Event) | 傳遞資料 (Payload) | 前端對應職責 |
|---|---|---|
status |
{ phase: "analyzing" / "generating" } |
切換主狀態文字與動畫 |
intent |
{ isMultiIntent: boolean, tags: [] } |
呈現意圖識別標籤 |
plan |
{ totalSteps: 3, actions: [...] } |
繪製 GOAP 思維鏈步驟清單 |
step |
{ actionName: "fetch", status: "DONE" } |
即時更新單一步驟勾選狀態與耗時 |
chunk |
string (JSON Spec 局部增量字串) |
累積至緩衝區並進行漸進式解析 |
complete |
{ executionCost: 0.002, durationMs: 1400 } |
流程結束,呈現最終成本與指標 |
error |
{ code: 500, message: "..." } |
捕捉例外並呈現友善錯誤面板 |
前端收到 chunk 後,必須經過三層防護管道:
jsonAccumulator += chunkData。lenientParseSpec):使用正規化與棧(Stack)演算法,自動為未閉合的引號補 ",為未閉合的大括號補 }。sanitizeSpec):
props: {}。children: []。children 包含尚未抵達的 Key,自動過濾以防渲染中斷。以下提供:
useDashboardStream)
// src/utils/lenientJsonParser.ts
import { DashboardSpec } from '../components/renderer/FlatTreeRenderer';
/**
* 寬鬆 JSON 解析器
* 自動嘗試修補串流中尚未閉合的引號與大括號
*/
export function lenientParseSpec(rawJson: string): any {
if (!rawJson || rawJson.trim() === '') return null;
let cleaned = rawJson.trim();
// 1. 嘗試直接解析
try {
return JSON.parse(cleaned);
} catch (e) {
// 進入修補模式
}
// 2. 修補未閉合引號
const quoteCount = (cleaned.match(/(?<!\\)"/g) || []).length;
if (quoteCount % 2 !== 0) {
cleaned += '"';
}
// 3. 修補未閉合的大括號與中括號
const openBraces = (cleaned.match(/\{/g) || []).length;
const closeBraces = (cleaned.match(/\}/g) || []).length;
for (let i = 0; i < openBraces - closeBraces; i++) {
cleaned += '}';
}
try {
return JSON.parse(cleaned);
} catch (e) {
return null; // 若修補後依然不合法,靜默等待下一批 Chunk
}
}
/**
* 規格安全淨化器
* 補齊缺失屬性並過濾不存在的子節點參照
*/
export function sanitizeSpec(parsed: any): DashboardSpec | null {
if (!parsed || typeof parsed !== 'object') return null;
if (!parsed.root || !parsed.elements || typeof parsed.elements !== 'object') return null;
const validElements: Record<string, any> = {};
const allKeys = new Set(Object.keys(parsed.elements));
for (const [key, element] of Object.entries<any>(parsed.elements)) {
if (!element || typeof element !== 'object' || !element.type) continue;
// 防禦 1: 補齊 props
const props = element.props && typeof element.props === 'object' ? element.props : {};
// 防禦 2: 過濾尚未抵達的子節點 key
const rawChildren = Array.isArray(element.children) ? element.children : [];
const safeChildren = rawChildren.filter((childKey: string) => allKeys.has(childKey));
validElements[key] = {
type: element.type,
props: props,
children: safeChildren
};
}
// 確保 root 節點存在
if (!validElements[parsed.root]) return null;
return {
root: parsed.root,
elements: validElements
};
}
useDashboardStream 自製 Hook// src/hooks/useDashboardStream.ts
import { useState, useRef } from 'react';
import { DashboardSpec } from '../components/renderer/FlatTreeRenderer';
import { lenientParseSpec, sanitizeSpec } from '../utils/lenientJsonParser';
export interface StreamProgress {
phase: string;
steps: Array<{ name: string; status: 'PENDING' | 'RUNNING' | 'DONE' }>;
costUsd: number;
}
export function useDashboardStream() {
const [spec, setSpec] = useState<DashboardSpec | null>(null);
const [progress, setProgress] = useState<StreamProgress>({ phase: 'IDLE', steps: [], costUsd: 0 });
const [isLoading, setIsLoading] = useState(false);
const jsonAccumulatorRef = useRef("");
const startStream = async (query: string) => {
setIsLoading(true);
setSpec(null);
jsonAccumulatorRef.current = "";
setProgress({ phase: 'CONNECTING', steps: [], costUsd: 0 });
try {
const response = await fetch('/api/dashboard/generate', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ query })
});
if (!response.body) throw new Error("瀏覽器不支援 ReadableStream");
const reader = response.body.getReader();
const decoder = new TextDecoder();
let buffer = "";
while (true) {
const { value, done } = await reader.read();
if (done) break;
buffer += decoder.decode(value, { stream: true }).replace(/\r\n/g, "\n");
const rawEvents = buffer.split("\n\n");
buffer = rawEvents.pop() || ""; // 留下半截待下輪拼接
for (const raw of rawEvents) {
const lines = raw.split("\n");
let eventType = "message";
let dataStr = "";
for (const line of lines) {
if (line.startsWith("event:")) eventType = line.replace("event:", "").trim();
if (line.startsWith("data:")) dataStr = line.replace("data:", "").trim();
}
if (!dataStr) continue;
// 處理不同類型的 SSE 事件
if (eventType === 'status') {
const data = JSON.parse(dataStr);
setProgress(prev => ({ ...prev, phase: data.phase }));
} else if (eventType === 'chunk') {
jsonAccumulatorRef.current += dataStr;
const parsed = lenientParseSpec(jsonAccumulatorRef.current);
const safe = sanitizeSpec(parsed);
if (safe) {
setSpec(safe); // 漸進式更新畫面!
}
} else if (eventType === 'complete') {
const data = JSON.parse(dataStr);
setProgress(prev => ({ ...prev, phase: 'COMPLETED', costUsd: data.executionCost }));
}
}
}
} catch (err: any) {
setProgress(prev => ({ ...prev, phase: 'ERROR' }));
console.error("串流中斷或發生錯誤:", err);
} finally {
setIsLoading(false);
}
};
return { spec, progress, isLoading, startStream };
}
package com.antechinus.travel.web;
import com.antechinus.travel.spec.DashboardSpec;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.springframework.http.MediaType;
import org.springframework.web.bind.annotation.*;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.io.IOException;
import java.util.Map;
import java.util.concurrent.CompletableFuture;
/**
* 儀表板生成串流端點
* 透過 POST 接收查詢並回傳 SSE 串流事件
*/
@RestController
@RequestMapping("/api/dashboard")
public class DashboardStreamController {
private final ObjectMapper objectMapper = new ObjectMapper();
public record GenerateRequest(String query) {}
@PostMapping(value = "/generate", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter streamDashboard(@RequestBody GenerateRequest request) {
// 設定 60 秒超時
SseEmitter emitter = new SseEmitter(60_000L);
CompletableFuture.runAsync(() -> {
try {
// 1. 發送分析狀態
sendEvent(emitter, "status", Map.of("phase", "ANALYZING"));
Thread.sleep(300);
// 2. 發送規劃思維鏈
sendEvent(emitter, "plan", Map.of(
"steps", List.of("讀取歷史旅程", "計算年度指標", "組裝儀表板規格")
));
Thread.sleep(400);
// 3. 模擬 Action 完成
sendEvent(emitter, "step", Map.of("actionName", "讀取歷史旅程", "status", "DONE"));
// 4. 模擬分批 Chunk 傳送 Spec
DashboardSpec sampleSpec = DashboardSpec.sample();
String fullJson = objectMapper.writeValueAsString(sampleSpec);
// 模擬將 JSON 拆為 3 個 Chunk 依序送出
int chunkSize = fullJson.length() / 3;
for (int i = 0; i < fullJson.length(); i += chunkSize) {
int end = Math.min(i + chunkSize, fullJson.length());
String chunk = fullJson.substring(i, end);
sendEvent(emitter, "chunk", chunk);
Thread.sleep(200); // 模擬生成間隔
}
// 5. 發送收尾事件
sendEvent(emitter, "complete", Map.of(
"executionCost", 0.0018,
"durationMs", 1250
));
emitter.complete();
} catch (Exception e) {
emitter.completeWithError(e);
}
});
return emitter;
}
private void sendEvent(SseEmitter emitter, String eventName, Object data) throws IOException {
String payload = (data instanceof String str) ? str : objectMapper.writeValueAsString(data);
emitter.send(SseEmitter.event().name(eventName).data(payload));
}
}
X-Accel-Buffering: no,或在 Nginx 配置 proxy_buffering off;,強制 Nginx 不得快取 SSE 封包。\r\n)導致事件解析破裂
\r\n,若用 split("\n\n") 會導致事件無法正確切割。.replace(/\r\n/g, "\n") 進行全域正規化。controller.abort()
useEffect 的 Cleanup 函式中呼叫 AbortController.abort() 主動中斷 fetch 連線。| 評估維度 | ❌ 瀏覽器原生 EventSource | ❌ WebSocket | ✅ Fetch + ReadableStream SSE |
|---|---|---|---|
| HTTP Method 支援 | 僅支援 GET,無法帶 Body | 獨立二進位/文字雙向協定 | 支援標準 POST 與 JSON Body |
| 協定複雜度 | 極低 | 較高(需處理心跳、握手、狀態維護) | 輕量標準 HTTP/2 串流 |
| 防火牆穿透性 | 良好 | 部分企業代理伺服器(Proxy)會攔截 | 100% 相容既有 HTTP 防火牆與 CDN |
| 半截 JSON 漸進渲染 | 支援(但難傳參數) | 支援 | 完美結合 lenientParse 與狀態漸進展示 |
這張截圖攝於查詢送出後、生成尚未完成的瞬間:前四個 GOAP action 已完成(右側各自帶耗時),「生成儀表板」正在執行、狀態列顯示「串流渲染中…」,左側儀表板已經先渲染出第一個元素(標題)——這就是 status → plan → step → chunk 事件管線與半截 JSON 容錯的實際樣貌。

useDashboardStream Hook,透過 fetch + ReadableStream 接收後端發送的 status 與 chunk 事件。{"root":"r1","elements":{"r1":{"type":"Stack"),驗證 lenientParseSpec 是否能自動補全大括號並成功由 sanitizeSpec 淨化。