
在昨天(Day 27)中,我們成功建立了一套強大的 TravelDashboardAgent,能夠透過 GOAP 演算法完成複雜的多步驟規劃與資料組裝。
然而,如果只是把 Agent 的最終結果打包成一個常規的 REST JSON 回傳,使用者端在按下按鈕後將會面對長達 4~8 秒的「黑盒等待期」——不知道系統是在查詢資料庫、調用 LLM 思考,還是已經卡死當機。
今天,我們要搭建前後端整合的第三塊核心基石——「SSE 進度事件與串流轉發管線」。
我們將解決一個在 Spring Boot + Agentic AI 整合中最高頻的架構難題:Embabel 的生命週期事件(PlanFormulated、ActionStart、ActionComplete)是屬於「平台層全域事件匯流排」,而前端 HTTP 請求則是「單一隔離的 SseEmitter」——我們該如何精確且安全地將全域事件路由到對應使用者的連線上?
EventListener 若沒有做好 Trace 隔離,使用者 A 會看到使用者 B 的思維步驟。onTimeout 與 onCompletion 清理快取,伺服器記憶體將迅速耗盡。assembleDashboard 是純 Java 毫秒級運算,沒有像 LLM 一樣逐字產出的 Token,如果不加處理,前端將無法觸發漸進渲染動畫。
SseEmitter 並向 SseProgressBroker 依據 TraceId 註冊。SseEmitter 生命週期治理SseEmitter 是 Spring MVC 原生支援非同步推播的核心物件。生產級別必須配置三個監聽回呼:
.onCompletion(() -> broker.remove(traceId)):正常完成時釋放資源。.onTimeout(() -> broker.remove(traceId)):超時關閉。.onError(e -> broker.remove(traceId)):網路異常中斷時清理。SseProgressBroker 負責將 Embabel 內部事件精準轉譯為前端定義好的 7 類標準 SSE 協定:
AgentProcessPlanFormulatedEvent $\rightarrow$ event: plan(推播整體規劃路徑)ActionExecutionStartEvent $\rightarrow$ event: step, data: { status: "RUNNING" }
ActionExecutionCompletedEvent $\rightarrow$ event: step, data: { status: "DONE" }
當純 Java Action 產出完整的 DashboardSpec 時,Controller 會透過演算法將其拆解為數個邏輯 Chunk(例如先送 Layout,再送 Cards,最後送 Table Rows),並以 50ms 間隔分批發送 event: chunk,在確保資料完整性的同時,賦予前端生動的漸進生長動畫。
以下提供:
SseProgressBroker)
DashboardStreamingController)
package com.antechinus.travel.streaming;
import com.embabel.agent.event.*;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;
import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
import java.io.IOException;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
/**
* SSE 進度事件轉發仲介器 (Active Connection Broker)
* 負責將 Embabel 全域生命週期事件精確路由至對應使用者的 SseEmitter 連線
*/
@Component
public class SseProgressBroker implements AgenticEventListener {
private static final Logger log = LoggerFactory.getLogger(SseProgressBroker.class);
private final ObjectMapper objectMapper = new ObjectMapper();
// 儲存進行中連線的執行緒安全容器 (TraceId -> SseEmitter)
private final ConcurrentHashMap<String, SseEmitter> activeEmitters = new ConcurrentHashMap<>();
/**
* 註冊使用者的 SseEmitter 連線
*/
public void register(String traceId, SseEmitter emitter) {
activeEmitters.put(traceId, emitter);
emitter.onCompletion(() -> unregister(traceId));
emitter.onTimeout(() -> unregister(traceId));
emitter.onError(e -> unregister(traceId));
}
/**
* 移除並清理連線
*/
public void unregister(String traceId) {
activeEmitters.remove(traceId);
log.debug("[SSE Broker] 已清理連線資源,TraceId: {}", traceId);
}
/**
* 監聽 Embabel 核心事件並轉換為前端 SSE 事件
*/
@Override
public void onProcessEvent(AgentProcessEvent event) {
SseEmitter emitter = activeEmitters.get(event.traceId());
if (emitter == null) {
return; // 若該 TraceId 無活躍連線則略過
}
try {
if (event instanceof PlanFormulatedEvent planEvent) {
send(emitter, "plan", Map.of(
"goal", planEvent.goalDescription(),
"steps", planEvent.actionNames(),
"totalSteps", planEvent.actionNames().size()
));
} else if (event instanceof ActionStartedEvent started) {
send(emitter, "step", Map.of(
"actionName", started.actionName(),
"status", "RUNNING"
));
} else if (event instanceof ActionCompletedEvent completed) {
send(emitter, "step", Map.of(
"actionName", completed.actionName(),
"status", "DONE",
"durationMs", completed.durationMs()
));
}
} catch (Exception e) {
log.error("[SSE Broker] 推播事件失敗,TraceId: {}", event.traceId(), e);
unregister(event.traceId());
}
}
/**
* 底層 SSE 事件發射輔助方法
*/
public void send(SseEmitter emitter, String eventName, Object payload) throws IOException {
String json = (payload instanceof String str) ? str : objectMapper.writeValueAsString(payload);
emitter.send(SseEmitter.event().name(eventName).data(json));
}
}
package com.antechinus.travel.web;
import com.antechinus.travel.domain.CustomerQuery;
import com.antechinus.travel.domain.UserQuery;
import com.antechinus.travel.spec.DashboardSpec;
import com.antechinus.travel.streaming.SseProgressBroker;
import com.embabel.agent.api.EmbabelClient;
import com.embabel.agent.api.ProcessResult;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
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.UUID;
import java.util.concurrent.CompletableFuture;
/**
* 儀表板動態生成串流端點
*/
@RestController
@RequestMapping("/api/dashboard")
@CrossOrigin(origins = "*") // 支援前端跨域開發
public class DashboardStreamController {
private static final Logger log = LoggerFactory.getLogger(DashboardStreamController.class);
private final EmbabelClient embabelClient;
private final SseProgressBroker progressBroker;
private final ObjectMapper objectMapper = new ObjectMapper();
public DashboardStreamController(EmbabelClient embabelClient, SseProgressBroker progressBroker) {
this.embabelClient = embabelClient;
this.progressBroker = progressBroker;
}
public record QueryRequest(String query) {}
@PostMapping(value = "/stream", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public SseEmitter streamDashboard(@RequestBody QueryRequest request) {
String traceId = "tr-" + UUID.randomUUID().toString().substring(0, 8);
SseEmitter emitter = new SseEmitter(60_000L); // 60 秒超時
// 1. 向 Broker 註冊連線
progressBroker.register(traceId, emitter);
// 2. 異步啟動 Agent 流程
CompletableFuture.runAsync(() -> {
try {
// 階段 A: 發送初始狀態
progressBroker.send(emitter, "status", Map.of("phase", "ANALYZING", "traceId", traceId));
// 階段 B: 執行 Embabel Agent 流程 (內部會自動觸發 EventListener)
ProcessResult<DashboardSpec> result = embabelClient.goal(DashboardSpec.class)
.withTraceId(traceId)
.withInput(new UserQuery(request.query()))
.run();
if (result.isSuccess()) {
DashboardSpec spec = result.getOutput();
String fullJson = objectMapper.writeValueAsString(spec);
// 階段 C: 合成 Chunk 分批推播 (Synthetic Chunking)
progressBroker.send(emitter, "status", Map.of("phase", "RENDERING"));
int chunkSize = Math.max(fullJson.length() / 4, 30);
for (int i = 0; i < fullJson.length(); i += chunkSize) {
int end = Math.min(i + chunkSize, fullJson.length());
progressBroker.send(emitter, "chunk", fullJson.substring(i, end));
Thread.sleep(60); // 營造流暢生長效果
}
// 階段 D: 發送完成事件
progressBroker.send(emitter, "complete", Map.of(
"executionCost", result.getCostUsd(),
"durationMs", result.getDurationMs(),
"status", "SUCCESS"
));
} else {
progressBroker.send(emitter, "error", Map.of(
"code", 500,
"message", "Agent 規劃失敗: " + result.getErrorMessage()
));
}
emitter.complete();
} catch (Exception e) {
log.error("[Controller] 串流執行例外,TraceId: {}", traceId, e);
try {
progressBroker.send(emitter, "error", Map.of("code", 500, "message", e.getMessage()));
} catch (IOException ignored) {}
emitter.completeWithError(e);
} finally {
progressBroker.unregister(traceId);
}
});
return emitter;
}
}
CompletableFuture.runAsync,導致 Spring Web 執行緒被長時間佔用,HTTP 請求遲遲無法開啟串流通道。SseEmitter 實例,所有 Agent 計算一律放入非同步執行緒池。finally 區塊中 unregister(traceId)
try-finally 確保在任何情況下都能徹底清理 TraceId。text/event-stream
@CrossOrigin(origins = "*") 或於 Spring Security 中統一配置 SSE Headers。| 評估維度 | ❌ 傳統同步 REST 回傳 (Bad) | ✅ SseProgressBroker 串流體系 (Good) |
|---|---|---|
| 等待體驗 | 盲等 8 秒,使用者容易以為系統卡死重複點擊 | 毫秒級接收 status 與 plan,進度條動態反饋 |
| 可觀測性 | 只有伺服器後台 Log 看得到進度 | 思維鏈與每一步 Action 執行細節即時展現給使用者 |
| 記憶體管理 | 無狀態維護,但並發高時線程易被塞滿 | 精確的 SseEmitter 生命週期回呼,超時自動回收 |
| UI 呈現 | 最後瞬間跳出整張 Dashboard | 透過 Chunk 漸進式生長,展現極致的 Generative UI 體驗 |
查詢執行中的瞬間:意圖判斷已完成(單一查詢)、Autonomy 正在選擇 Agent,GOAP 規劃出的 5 個 action 以 plan 事件先行送達前端排成清單,第 1 步 extractChurnParams 標示「執行中…」——使用者從送出查詢的第一秒起,就能看見 agent 正在規劃與執行什麼。

SseProgressBroker,實作 AgenticEventListener 介面並完成 TraceId 路由。DashboardStreamController 實作 /api/dashboard/stream,並使用 curl -N -X POST http://localhost:8080/api/dashboard/stream -H "Content-Type: application/json" -d "{\"query\":\"查客戶1001\"}" 驗證串流輸出。