iT邦幫忙

2026 iThome 鐵人賽

DAY 28
0

https://ithelp.ithome.com.tw/upload/images/20260828/2016129022m55iAQdB.png

讓使用者看見 agent 正在規劃與執行。

在昨天(Day 27)中,我們成功建立了一套強大的 TravelDashboardAgent,能夠透過 GOAP 演算法完成複雜的多步驟規劃與資料組裝。

然而,如果只是把 Agent 的最終結果打包成一個常規的 REST JSON 回傳,使用者端在按下按鈕後將會面對長達 4~8 秒的「黑盒等待期」——不知道系統是在查詢資料庫、調用 LLM 思考,還是已經卡死當機

今天,我們要搭建前後端整合的第三塊核心基石——「SSE 進度事件與串流轉發管線」

我們將解決一個在 Spring Boot + Agentic AI 整合中最高頻的架構難題:Embabel 的生命週期事件(PlanFormulated、ActionStart、ActionComplete)是屬於「平台層全域事件匯流排」,而前端 HTTP 請求則是「單一隔離的 SseEmitter」——我們該如何精確且安全地將全域事件路由到對應使用者的連線上?


1. 今天要解決的痛點與核心觀念

痛點背景:Agent 串流管線的三大工程雷區

  1. 事件交叉污染(Event Crosstalk):若多個使用者同時發起查詢,全域的 EventListener 若沒有做好 Trace 隔離,使用者 A 會看到使用者 B 的思維步驟。
  2. 連線記憶體洩漏(Emitter Leaks):當前端關閉瀏覽器或連線逾時時,後端若未在 onTimeoutonCompletion 清理快取,伺服器記憶體將迅速耗盡。
  3. 純 Java Action 無 Chunk 可吐:因為後端 assembleDashboard 是純 Java 毫秒級運算,沒有像 LLM 一樣逐字產出的 Token,如果不加處理,前端將無法觸發漸進渲染動畫。

觀念圖解:Active Connection Broker 架構模式

https://ithelp.ithome.com.tw/upload/images/20260828/201612904icVgaJ2VV.png

  1. 連線註冊 (Register):Controller 建立 SseEmitter 並向 SseProgressBroker 依據 TraceId 註冊。
  2. 全域事件路由 (Route by TraceId):Embabel 引擎發布生命週期事件時,Broker 精確轉發至對應客戶端的 SSE 通道。
  3. 安全釋放 (Unregister):連線結束或異常時自動清理 Map 快取,杜絕記憶體洩漏。

2. 官方核心技術依據與架構深度

1. Spring MVC SseEmitter 生命週期治理

SseEmitter 是 Spring MVC 原生支援非同步推播的核心物件。生產級別必須配置三個監聽回呼:

  • .onCompletion(() -> broker.remove(traceId)):正常完成時釋放資源。
  • .onTimeout(() -> broker.remove(traceId)):超時關閉。
  • .onError(e -> broker.remove(traceId)):網路異常中斷時清理。

2. 進度事件轉換層(Event Translation Layer)

SseProgressBroker 負責將 Embabel 內部事件精準轉譯為前端定義好的 7 類標準 SSE 協定:

  • AgentProcessPlanFormulatedEvent $\rightarrow$ event: plan(推播整體規劃路徑)
  • ActionExecutionStartEvent $\rightarrow$ event: step, data: { status: "RUNNING" }
  • ActionExecutionCompletedEvent $\rightarrow$ event: step, data: { status: "DONE" }

3. 合成串流分塊(Synthetic Chunking Strategy)

當純 Java Action 產出完整的 DashboardSpec 時,Controller 會透過演算法將其拆解為數個邏輯 Chunk(例如先送 Layout,再送 Cards,最後送 Table Rows),並以 50ms 間隔分批發送 event: chunk,在確保資料完整性的同時,賦予前端生動的漸進生長動畫。


3. 完整程式碼實戰(Production-Ready Code)

以下提供:

  1. 執行緒安全的連線仲介轉發器(SseProgressBroker
  2. 生產級串流發射端點(DashboardStreamingController

1. 全域事件轉發器實作(SseProgressBroker)

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));
    }
}

2. 生產級 SSE 串流 Controller 實作

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;
    }
}

4. 生產環境避坑指南與對比分析

常見踩雷與除錯秘訣

  1. 雷區一:在 Controller 主執行緒中阻塞執行 Agent
    • 現象:沒有使用 CompletableFuture.runAsync,導致 Spring Web 執行緒被長時間佔用,HTTP 請求遲遲無法開啟串流通道。
    • 解法:Controller 必須立即回傳 SseEmitter 實例,所有 Agent 計算一律放入非同步執行緒池。
  2. 雷區二:忘記在 finally 區塊中 unregister(traceId)
    • 現象:拋出例外時,Emitter 殘留在 Map 容器中,導致記憶體洩漏與垃圾收集器(GC)負擔加重。
    • 解法:務必使用 try-finally 確保在任何情況下都能徹底清理 TraceId。
  3. 雷區三:跨網域 CORS 配置未開放 text/event-stream
    • 現象:前端本機開發(localhost:5173)請求後端(localhost:8080)時,被瀏覽器 CORS 策略攔截。
    • 解法:在 Controller 上標記 @CrossOrigin(origins = "*") 或於 Spring Security 中統一配置 SSE Headers。

串流架構 Good vs Bad 對比表

評估維度 ❌ 傳統同步 REST 回傳 (Bad) ✅ SseProgressBroker 串流體系 (Good)
等待體驗 盲等 8 秒,使用者容易以為系統卡死重複點擊 毫秒級接收 statusplan,進度條動態反饋
可觀測性 只有伺服器後台 Log 看得到進度 思維鏈與每一步 Action 執行細節即時展現給使用者
記憶體管理 無狀態維護,但並發高時線程易被塞滿 精確的 SseEmitter 生命週期回呼,超時自動回收
UI 呈現 最後瞬間跳出整張 Dashboard 透過 Chunk 漸進式生長,展現極致的 Generative UI 體驗

5. 實機畫面:進度事件即時推到前端

查詢執行中的瞬間:意圖判斷已完成(單一查詢)、Autonomy 正在選擇 Agent,GOAP 規劃出的 5 個 action 以 plan 事件先行送達前端排成清單,第 1 步 extractChurnParams 標示「執行中…」——使用者從送出查詢的第一秒起,就能看見 agent 正在規劃與執行什麼。

https://ithelp.ithome.com.tw/upload/images/20260828/20161290WyaLEqhOte.png


6. 今日動手實作任務與發文備註

🛠️ 今日實作任務

  1. 實作 SseProgressBroker:在專案中建立 SseProgressBroker,實作 AgenticEventListener 介面並完成 TraceId 路由。
  2. 建立 SSE 串流端點:在 DashboardStreamController 實作 /api/dashboard/stream,並使用 curl -N -X POST http://localhost:8080/api/dashboard/stream -H "Content-Type: application/json" -d "{\"query\":\"查客戶1001\"}" 驗證串流輸出。
  3. 思考題:為什麼在全域監聽器中透過 TraceId 來做連線路由,比在每個 Action 內部直接注入 SseEmitter 更符合軟體工程的「關注點分離」原則?

上一篇
Day 27:把資料整理成畫面
下一篇
Day 29:讓畫面慢慢長出來
系列文
讓 AI Agent 真的做事:用 Embabel 打造可控、可測試的智慧 Dashboard29
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言