Langgraph4j - 基础介绍和案例演示(一)

Langgraph4j 介绍

它是什么?

从根本上说,LangGraph4j 的本质是一个有状态、可循环的图式流程编排引擎。它是 Python LangGraph 的 Java 实现,专门用来解决 复杂 Agent 系统中 “怎么走流程、怎么管状态、怎么协作” 的问题。它与 LangChain4j 的不同,可以用一句话类比:

  • LangChain4j 是 “零件箱 + 接口适配层”,负责把大模型、Embedding、Tool、RAG 这些能力接入 Java 应用(像高速公路,从入口到出口基本是一条直线)。
  • LangGraph4j 是 “搭积木的图纸 + 流程引擎”,负责用 StateGraph(状态图)Node(节点)Edge(边) 把多个 Agent/步骤组织成有分支、有循环、有共享状态、可中断恢复的工作流(像城市路网,可以分叉、掉头、并行)。


为什么需要它?

你也许会问:“既然 LangChain4j 也能处理相关问题,为什么还需要 LangGraph4j?” LangChain4j 确实能处理单轮/多轮工具调用、简单 ReAct Agent、线性 Chain、基础 RAG 等场景,它的执行模型本质上是有向无环图(DAG)/ 链式顺序执行,状态管理主要靠 ChatMemory且多是无状态的单次工作流视角。但当你的 Agent 需求往前再走一步,LangChain4j 就会碰到架构级的天花板,而这些正是 LangGraph4j 存在的意义:

  • 循环与自我修正:
    • LangChain4j 的局限:链式结构很难原生表达 “思考 → 行动 → 观察 → 不满意 → 回到上一步重新思考” 这种环。虽然能用 Loop或递归硬编码,但复杂度一上来就难以维护。
    • LangGraph4j 的能力:原生支持循环边,图结构允许节点执行完回到起点或其他上游节点,天然适配 Agent 的“反思、重试、多轮谈判、迭代生成”等场景。
  • 显式且持久的全局状态管理:
    • LangChain4j 的局限:状态分散在 ChatMemory 或链式传递中,没有统一的、版本化的、可持久化的共享状态对象,跨长时间(甚至跨天)的恢复需要自己手写存储逻辑。
    • LangGraph4j 的能力:强制定义 AgentState 作为全局共享状态,每步执行自动更新,并内置 Checkpoint(检查点) 机制,可存到 Redis/PG 中,系统崩溃、人工介入后能从断点精确恢复(Time Travel)。
  • 复杂条件路由与 LLM 自主决策:
    • LangChain4j 的局限:分支大多是开发者硬编码的 if-else / Router,LLM 不能真正动态决定“退回哪一步”或“在多分支中自主选路”。
    • LangGraph4j 的能力:提供 Conditional Edge(条件边),路由函数可以由 LLM 输出驱动,让模型根据当前 State 自主决定下一步走哪个节点,实现真正的动态规划。
  • 多 Agent 协作与并行编排:
    • LangChain4j 的局限:能做多实例调用,但缺乏原生的 Supervisor + Worker、Fan-out 并行聚合、子图嵌套 等复杂多 Agent 拓扑编排能力。
    • LangGraph4j 的能力:天生为 多智能体协作 设计,支持并行节点、子图、共享状态读写,适合构建分工明确的 Agent 团队(如:分诊 Agent → 检索 Agent → 审核 Agent → 汇总 Agent)。
  • 人机协同:
    • LangChain4j 的局限:Tool 执行完就继续,没有原生 “在执行途中暂停、等人工审批/修改后再继续” 的机制。
    • LangGraph4j 的能力:通过 interrupt 等机制在任意节点暂停图执行,等待人工输入后再恢复,这对医疗、金融等需要人工把关的企业级场景至关重要。

LangGraph4j 和 LangChain4j 不是替代,而是不同层级的互补关系,在实际项目中通常会一起用。简单说就是:LangChain4j 解决 “每个节点里用什么能力干活”,LangGraph4j 解决 “节点之间怎么连、状态怎么传、流程怎么走” 的问题。

  • LangChain4j(能力层):负责模型调用(Ollama / OpenAI 等)、Embedding、向量检索 RAG、Tool 定义与执行、Prompt 模板、流式输出。
  • LangGraph4j(编排层):负责定义 StateGraph、节点逻辑(内部可调用 LangChain4j 的 Agent/Tool)、条件路由、循环控制、状态持久化、多 Agent 协调、断点恢复。


Langgraph4j 案例

这里我们以一个常见的 “客服工单智能路由与自动处理” 为例,使用 LangGraph4j 1.8.20、LangChain4j 1.17.2 以及 Spring Boot WebFlux 进行案例演示说明。


业务场景

当用户输入一个诉求(例如:“我的账号被锁定了” 或 “你们这个产品今天打折吗”),系统会经过以下图(Graph)流程:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
          ┌───────────────┐
│ START │
└───────┬───────┘


┌───────────────────┐
│ 1. IntentAnalyze │ (意图分析:分类为【技术/业务/其它】)
└───────┬───────────┘


┌───────────────────┐
│ 2. RouterNode │ (路由决策:条件边分支)
└─┬───────────────┬─┘
│ │
[技术问题]│ │[业务咨询]
▼ ▼
┌──────────────────┐ ┌──────────────────┐
│ 3A. TechSupport │ │ 3B. SalesSupport │ (分别调用 LLM 专属 Prompt)
└─────────┬────────┘ └─────────┬────────┘
│ │
└───────────┬───────────┘


┌───────────────────┐
│ 4. QualityGuard │ (质量检测:安全或合规检查)
└─────────┬─────────┘


/─────────────────\
< Is Approved? > (条件边:是否允许输出?)
\─────────────────/
/ \
[是] / \ [否:需人工接入或重构]
/ \
▼ ▼
┌───────────────┐ ┌──────────────────┐
│ END │ │ 5. HumanReview │ (人工接管/状态挂起)
└───────────────┘ └──────────────────┘


依赖配置

父项目配置,请参考本站 Langchain4j - 基础工程的构建以及两套API测试案例 - 父项目-POM

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-webflux</artifactId>
</dependency>

<dependency>
<groupId>org.bsc.langgraph4j</groupId>
<artifactId>langgraph4j-core</artifactId>
</dependency>

<dependency>
<groupId>dev.langchain4j</groupId>
<artifactId>langchain4j-core</artifactId>
</dependency>

<!-- 这里以 OpenAI 适配器为例,你也可以换成 Ollama 或 DeepSeek -->
<dependency>
<groupId>dev.langchain4j</groupId>
<artifactId>langchain4j-open-ai</artifactId>
</dependency>

<dependency>
<groupId>ch.qos.logback</groupId>
<artifactId>logback-classic</artifactId>
</dependency>

<dependency>
<groupId>com.fasterxml.jackson.core</groupId>
<artifactId>jackson-databind</artifactId>
</dependency>


状态定义与通道机制

定义 AgentState

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
import org.bsc.langgraph4j.state.AgentState;
import org.bsc.langgraph4j.state.Channel;
import org.bsc.langgraph4j.state.Channels;
import java.util.ArrayList;
import java.util.Map;

/**
* 定义 SupportState 的根本目的,就是为了用声明式的方式,提前告诉 LangGraph4j 引擎:
* 我的工作流里有哪些变量、它们默认值是多少、以及当多个节点改写它们时,应该采取覆盖还是追加策略。
*
* 为什么继承 AgentState?
* AgentState 在底层实际上是一个不可变的(Immutable)数据信封。
* 每个节点执行完毕后,只能返回一个包含新数据的微型 Map(更新增量)。框架收到这个 Map 后,会根据你在 SCHEMA 里定义的
* 规则,跟旧的状态合并,融合成一个新的 SupportState 传递给下一个节点。这种设计保证了状态在异步、并发环境下的线程安全。
*/
public class SupportState extends AgentState {

// 1. 定义 Key 常量
public static final String USER_QUERY = "userQuery";
public static final String INTENT = "intent";
public static final String AGENT_RESPONSE = "agentResponse";
public static final String IS_SAFE = "isSafe";
public static final String HISTORY = "history";

// 2. 定义 State Schema:描述每个 key 的合并/覆写规则
public static final Map<String, Channel<?>> SCHEMA = Map.of(
// Channels.base:收到新值时,直接丢弃旧值,保留新值。(适合 userQuery、intent 这种每次都在变化的单一变量)
USER_QUERY, Channels.base(() -> ""),
INTENT, Channels.base(() -> ""),
AGENT_RESPONSE, Channels.base(() -> ""),
IS_SAFE, Channels.base(() -> true),
// Channels.appender:收到新数据时,不覆盖旧数据,而是自动调用类似 list.add(newValue) 的操作,将数据追加到现有的 List 后面。(适合 history 这样的对话历史)
// 只要有新数据来,不要覆盖,请把新数据追加到现有的列表末尾;ArrayList::new 是说如果图刚启动时历史记录是空的,默认新建一个 ArrayList 准备着。
HISTORY, Channels.appender(ArrayList::new)
);

public SupportState(Map<String, Object> initData) {
super(initData);
}

// 辅助 Getter 方法,方便在节点中强类型读取
public String getUserQuery() {
return (String) value(USER_QUERY).orElse("");
}

public String getIntent() {
return (String) value(INTENT).orElse("UNKNOWN");
}

public String getAgentResponse() {
return (String) value(AGENT_RESPONSE).orElse("");
}

public Boolean isSafe() {
return (Boolean) value(IS_SAFE).orElse(true);
}
}


必要说明

在 LangGraph4j 中,状态的每一个 Key 都不再是 Map 里一个随意的 K-V 对,而被具象化为了一个通道(Channel)。以前的普通 Map 就像是一个公共储物柜,任何人都可以把里面的东西扔掉,换成自己的东西,或者把别人物品弄乱,没有任何秩序。Channel 管理的状态就像是一组有特定管道阀门的水管(Channel),每个水管(Channel)不仅负责存水,还自带规则(Reducer)和过滤器。

  • 它是 “有型(Typed)” 的:每一个通道在定义时就绑定了数据类型和初始逻辑。例如:Channels.appender(ArrayList::new) ,它明确知道自己管理的是一个 List。一旦定义,这个通道的数据结构就是确定的,框架会在底层处理好类型的装箱和拆箱。
  • 它拥有 “行为规则”(Reducer):这是 Channel 最核心的灵魂。当你向一个通道写入数据时,你不是在执行“覆盖” 操作,而是在向通道发送一个 “更新信号”。通道会根据自己内置的 Reducer(合并器)来决定如何融合新旧数据。

因为状态是由 Channel 管理的,LangGraph4j 的节点(Node)在执行完后,不需要(也不能)去手动修改整个 State。节点只需要返回一个极简的、只包含当前节点产出数据的增量 Map:

1
2
// 在某个 Node 中,我只需要返回这个:
return Map.of("intent", "TECH_SUPPORT");

当这个增量 Map 被提交给 LangGraph4j 引擎时,引擎会像下面这样运作:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
[ 节点 A 执行完毕 ] 

▼ 产出增量数据: { "intent": "TECH", "history": "Hello" }

┌──────┴────────────────────────────────────────┐
│ LangGraph 引擎分配调度 │
├──────────────────────┬────────────────────────┤
│ 对 "intent" 键: │ 对 "history" 键: │
│ 走向 base 通道 │ 走向 appender 通道 │
│ ↓ │ ↓ │
│ 【直接覆盖旧值】 │ 【调用 Reducer 追加到列表】│
└──────────────────────┴─────────────────────────┘

▼ 融合成全新的: SupportState

[ 传递给下一个节点 ]

# 这种设计被称为 State Redux(状态折叠/归约)。


定义工作节点

每个节点实现 NodeAction 接口,接收当前 State,返回更新后的 Map 数据。

IntentAnalyzeNode:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
import com.owlias.demo06.state.SupportState;
import dev.langchain4j.model.chat.ChatModel;
import org.bsc.langgraph4j.action.NodeAction;
import java.util.Map;

/**
* 节点 1:意图分析节点
*/
public class IntentAnalyzeNode implements NodeAction<SupportState> {

private final ChatModel chatModel;
public IntentAnalyzeNode(ChatModel chatModel) {
this.chatModel = chatModel;
}

@Override
public Map<String, Object> apply(SupportState state) {
String query = state.getUserQuery();
String prompt = """
分析以下用户的咨询意图。只输出以下三个标签之一:[TECH]、[SALES]、[OTHER]。
不要输出任何多余的解释。
用户咨询: "%s"
""".formatted(query);
String intent = chatModel.chat(prompt).trim();

// 返回需要更新的状态片断
return Map.of(SupportState.INTENT, intent);
}
}

TechSupportNode:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
import com.owlias.demo06.state.SupportState;
import dev.langchain4j.model.chat.ChatModel;
import org.bsc.langgraph4j.action.NodeAction;
import java.util.Map;

/**
* 节点 3A:技术支持专家节点
*/
public class TechSupportNode implements NodeAction<SupportState> {
private final ChatModel chatModel;
public TechSupportNode(ChatModel chatModel) {
this.chatModel = chatModel;
}

@Override
public Map<String, Object> apply(SupportState state) {
String prompt = "你是一位技术支持专家。请专业回答用户关于技术、系统、环境等问题(每次回答控制100字以内)。用户问题:"
+ state.getUserQuery();
String response = chatModel.chat(prompt);
return Map.of(SupportState.AGENT_RESPONSE, response);
}
}

SalesSupportNode:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
import com.owlias.demo06.state.SupportState;
import dev.langchain4j.model.chat.ChatModel;
import org.bsc.langgraph4j.action.NodeAction;
import java.util.Map;

/**
* 节点 3B:业务销售专家节点
*/
public class SalesSupportNode implements NodeAction<SupportState> {
private final ChatModel chatModel;

public SalesSupportNode(ChatModel chatModel) {
this.chatModel = chatModel;
}

@Override
public Map<String, Object> apply(SupportState state) {
String prompt = "你是一位金牌销售顾问。请简短介绍我们的产品(每次回答100字以内),解答价格、折扣和方案咨询。用户问题:"
+ state.getUserQuery();
String response = chatModel.chat(prompt);
return Map.of(SupportState.AGENT_RESPONSE, response);
}
}

QualityGuardNode:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
import com.owlias.demo06.state.SupportState;
import org.bsc.langgraph4j.action.NodeAction;
import java.util.Map;

/**
* 节点 4:内容合规性卫士
*/
public class QualityGuardNode implements NodeAction<SupportState> {
@Override
public Map<String, Object> apply(SupportState state) {
String response = state.getAgentResponse();

// 生产环境通常会接入敏感词过滤。这里做模拟演示:
// 假设 AI 提到了敏感词,或者回复为空,则判定不合规
boolean isSafe = !response.contains("密码") && !response.isEmpty();
return Map.of(SupportState.IS_SAFE, isSafe);
}
}


构建图与决策路由

现在我们利用 StateGraph 把所有的 Node 和 Edge 串联在一起,并定义条件边(Conditional Edges):

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
import com.owlias.demo06.node.IntentAnalyzeNode;
import com.owlias.demo06.node.QualityGuardNode;
import com.owlias.demo06.node.SalesSupportNode;
import com.owlias.demo06.node.TechSupportNode;
import com.owlias.demo06.state.SupportState;
import dev.langchain4j.model.chat.ChatModel;
import dev.langchain4j.model.openai.OpenAiChatModel;
import org.bsc.langgraph4j.CompiledGraph;
import org.bsc.langgraph4j.StateGraph;
import org.bsc.langgraph4j.action.AsyncEdgeAction;
import org.bsc.langgraph4j.action.AsyncNodeAction;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Primary;
import java.time.Duration;
import java.util.Map;

@Configuration
public class AgentGraphConfig {

@Primary
@Bean("localChatModel")
public ChatModel localChatModel() {
return OpenAiChatModel.builder() // 👈🏻
.baseUrl("http://localhost:11434/v1")
.apiKey("api_key_xxx")
.modelName("gemma3:1b")
.timeout(Duration.ofSeconds(300))
.maxRetries(3)
.logRequests(true)
.logResponses(true)
.build();
}

@Bean
public CompiledGraph<SupportState> customerSupportGraph(ChatModel model) throws Exception {

// 1. 初始化有状态图,绑定 State Schema 和构造器
StateGraph<SupportState> graph = new StateGraph<>(SupportState.SCHEMA, SupportState::new);

// 2. 注入节点 (推荐用 node_async 包装以获取最佳并发性能)
graph.addNode("intent_analyze", AsyncNodeAction.node_async(new IntentAnalyzeNode(model)));
graph.addNode("tech_support", AsyncNodeAction.node_async(new TechSupportNode(model)));
graph.addNode("sales_support", AsyncNodeAction.node_async(new SalesSupportNode(model)));
graph.addNode("quality_guard", AsyncNodeAction.node_async(new QualityGuardNode()));

// 3. 配置连线:入口点指向意图分析
graph.addEdge(StateGraph.START, "intent_analyze");

// 4. 条件路由分支:
graph.addConditionalEdges("intent_analyze",
// 根据意图路由到技术还是销售
AsyncEdgeAction.edge_async(state -> {
String intent = state.getIntent();
if (intent.contains("TECH")) {
return "route_to_tech";
} else {
return "route_to_sales";
}
}),
// 明确逻辑返回值与节点 ID 的映射关系
Map.of(
"route_to_tech", "tech_support",
"route_to_sales", "sales_support"
)
);

// 5. 业务处理完成后,强制汇总流向质量卫士节点
graph.addEdge("tech_support", "quality_guard");
graph.addEdge("sales_support", "quality_guard");

// 6. 第二个条件分支:如果安全则结束,不安全则走向人工审批/阻断
graph.addConditionalEdges("quality_guard",
AsyncEdgeAction.edge_async(state -> {
if (state.isSafe()) {
return "safe";
} else {
return "unsafe";
}
}),
Map.of(
"safe", StateGraph.END,
"unsafe", StateGraph.END // 这里简单处理,不安全也走向结束(实际生产可导向人工干预节点)
)
);

// 7. 编译生成可执行的 CompiledGraph
CompiledGraph<SupportState> compiledGraph = graph.compile();
System.out.println("\n" + compiledGraph.getGraph(GraphRepresentation.Type.MERMAID) + "\n"); // 可视化调试,也支持 PLANTUML 格式 👈🏻
return compiledGraph;
}
}


流式响应控制器

在 WebFlux 异步非阻塞架构下,我们可以通过 Flux 直接向前端推送图运行的中间状态:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
@RestController
public class SupportAgentController {

@Resource
private CompiledGraph<SupportState> compiledGraph;

/**
* 以 Server-Sent Events (SSE) 方式,流式输出工作流每个节点的运行状态
*/
@GetMapping(value = "/api/agent/ask", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<String> askAgent(@RequestParam String query) {
// 1. 初始化启动状态
Map<String, Object> inputs = Map.of(SupportState.USER_QUERY, query);
try {
// 2. 借助 LangGraph4j 的 stream 机制,并用 Reactor 的 Flux 包装
return Flux.fromIterable(compiledGraph.stream(inputs))
.map(nodeOutput -> {
String nodeName = nodeOutput.node(); // 当前执行完的节点名
SupportState stateSnapshot = nodeOutput.state(); // 当前状态快照

return """
Node: [%s]
- Intent: %s
- Safe: %s
- Current Response: %s
-------------------------
""".formatted(
nodeName,
stateSnapshot.getIntent(),
stateSnapshot.isSafe(),
stateSnapshot.getAgentResponse());
})
.onErrorReturn("Execution error occurred during agent processing.");
} catch (Exception e) {
e.printStackTrace();
return Flux.just("Error initializing Graph: " + e.getMessage());
}
}
}

测试验证:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
$ curl -N -G 'http://localhost:8080/api/agent/ask' \
--data-urlencode 'query=请问你们的企业版怎么收费,有折扣吗?'

data:Node: [__START__]
data:- Intent:
data:- Safe: true
data:- Current Response:
data:-------------------------
data:

data:Node: [intent_analyze]
data:- Intent: [SALES]
data:- Safe: true
data:- Current Response:
data:-------------------------
data:

data:Node: [sales_support]
data:- Intent: [SALES]
data:- Safe: true
data:- Current Response: 您好!很高兴为您服务。
data:
data:我们的企业版主要提供**高级数据分析和报告生成服务**,能帮助您更好地了解市场趋势、客户行为和运营效率。
data:
data:**收费方式:** 我们可以根据您的需求定制方案,采用按月或按年计费。初期,我们会提供免费试用,了解您的数据需求。
data:
data:**折扣:** 我们提供**分期付款和长期合作优惠**,以延长您的投资周期及降低成本。
data:
data:您想了解哪方面的具体数据分析? 或者您对价格和方案有什么疑问?请告诉我,我来详细解答您的问题!
data:-------------------------
data:

data:Node: [quality_guard]
data:- Intent: [SALES]
data:- Safe: true
data:- Current Response: 您好!很高兴为您服务。
data:
data:我们的企业版主要提供**高级数据分析和报告生成服务**,能帮助您更好地了解市场趋势、客户行为和运营效率。
data:
data:**收费方式:** 我们可以根据您的需求定制方案,采用按月或按年计费。初期,我们会提供免费试用,了解您的数据需求。
data:
data:**折扣:** 我们提供**分期付款和长期合作优惠**,以延长您的投资周期及降低成本。
data:
data:您想了解哪方面的具体数据分析? 或者您对价格和方案有什么疑问?请告诉我,我来详细解答您的问题!
data:-------------------------
data:

data:Node: [__END__]
data:- Intent: [SALES]
data:- Safe: true
data:- Current Response: 您好!很高兴为您服务。
data:
data:我们的企业版主要提供**高级数据分析和报告生成服务**,能帮助您更好地了解市场趋势、客户行为和运营效率。
data:
data:**收费方式:** 我们可以根据您的需求定制方案,采用按月或按年计费。初期,我们会提供免费试用,了解您的数据需求。
data:
data:**折扣:** 我们提供**分期付款和长期合作优惠**,以延长您的投资周期及降低成本。
data:
data:您想了解哪方面的具体数据分析? 或者您对价格和方案有什么疑问?请告诉我,我来详细解答您的问题!
data:-------------------------
data:

我们发现,上述执行的过程出现了:

1
[用户请求] -> 执行节点 A -> (等待 LLM 完备响应) -> 产生 StateA -> 执行节点 B -> (等待 LLM 完备响应) -> 产生 StateB -> [结束]


如何做到完全的流式响应

上述演示,虽然输出的整个过程是流式响应式的,但是 [TECH] 或 [SALES] 节点的执行,结果的吐出却是一次性的。我们需要区分两个不同维度的“流”:

  • 图的状态流 (Graph State Stream):当你在 Controller 里调用 graph.stream(inputs) 时,LangGraph4j 吐出的事件颗粒度是节点级别的(例如:“节点 A 执行完毕,状态已更新” -> “节点 B 开始执行”)。因为节点 B 必须等待节点 A 的输出作为上下文,所以图的骨架流转在逻辑上必须是前因后果的串行。
  • 令牌流 (Token Stream / SSE):这才是我们前端想要的“打字机效果”。当 TechSupportNode 正在执行时,大模型产生一字一句的文本,需要通过另一个通道实时推给前端。

如果每个节点的结果是一次性蹦出来的,说明节点在等待大模型完全生成完毕后,才把整个字符串放回 SupportState并结束节点。想要实现整个过程都是流式响应,要在维持图的复杂状态流转的同时,让用户看到打字机效果,比较标准的做法是:将大模型的 Token 流(StreamingChatModel)直接定向输出给 WebFlux 的响应流(或 SseEmitter),而图的状态只负责承载最终沉淀的数据。我们以 [TECH] 节点为例,进行改造介绍。


同时引入流式 LLM

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
@Primary
@Bean("localChatModel")
public ChatModel localChatModel() {
return OpenAiChatModel.builder() // 👈🏻
.baseUrl("http://localhost:11434/v1")
.apiKey("api_key_xxx")
.modelName("gemma3:1b")
.timeout(Duration.ofSeconds(300))
.maxRetries(3)
.logRequests(true)
.logResponses(true)
.build();
}

@Bean
public StreamingChatModel localStreamingChatModel() {
return OpenAiStreamingChatModel.builder() // 👈🏻
.baseUrl("http://localhost:11434/v1")
.apiKey("api_key_xxx")
.modelName("gemma3:1b")
.timeout(Duration.ofSeconds(300))
.logRequests(true)
.logResponses(true)
.build();
}


业务节点的改造

在负责生成长文本的节点(如 TechSupportNode)中,我们使用 StreamingChatModel。它可以在运行期间直接往当前请求的响应流里源源不断地塞文本:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
import com.owlias.demo06.state.SupportState;
import dev.langchain4j.model.chat.StreamingChatModel;
import dev.langchain4j.model.chat.response.ChatResponse;
import dev.langchain4j.model.chat.response.StreamingChatResponseHandler;
import org.bsc.langgraph4j.RunnableConfig;
import reactor.core.publisher.FluxSink;
import java.util.Map;
import java.util.Optional;
import java.util.concurrent.CompletableFuture;

/**
* 节点 3A:技术支持专家节点
*/
public class TechSupportNode { // 这里不再继承 NodeAction 或 AsyncNodeAction
private final StreamingChatModel streamingChatModel;

public TechSupportNode(StreamingChatModel streamingChatModel) {
this.streamingChatModel = streamingChatModel;
}

/**
* 节点执行方法
*/
public CompletableFuture<Map<String, Object>> execute(SupportState state, RunnableConfig config) {
// 👉🏻 健壮性:安全地获取 sseSink,避免 get() 报错
Optional<Object> sseSinkOpt = config.metadata("sseSink");
if (sseSinkOpt.isEmpty()) {
// 如果没有 sseSink,做非流式降级或更友好的错误提示,而不是直接崩掉整个流
return CompletableFuture.completedFuture(Map.of(
SupportState.AGENT_RESPONSE, "系统正忙,暂无法生成流式回答,请稍后再试。",
SupportState.IS_SAFE, true
));
}

// 从 RunnableConfig 中捞出瞬时响应流。其中 FluxSink<String>(或者自定义的事件对象)可以通过 LangGraph 的 RunnableConfig(配置/上下文)或 State 传进来。
@SuppressWarnings("unchecked")
FluxSink<String> sseSink = (FluxSink<String>) sseSinkOpt.get();

String userQuery = state.getUserQuery();
String prompt = "你是一位技术支持专家。请专业回答用户关于技术、系统、环境等问题(每次回答控制100字以内)。用户问题:" + userQuery;

StringBuilder fullResponse = new StringBuilder();
CompletableFuture<Map<String, Object>> future = new CompletableFuture<>();
streamingChatModel.chat(prompt, new StreamingChatResponseHandler() {
@Override
public void onPartialResponse(String partialResponse) {
// 👉🏻 再次防御:防止向一个已经关闭的 sink 发送数据报错
try {
// 只有在 sink 未被取消时才推送数据
if (!sseSink.isCancelled()) {
sseSink.next(partialResponse);
}
} catch (Exception e) {
// 打印日志:通常是前端主动断开了连接
}
fullResponse.append(partialResponse);
}

@Override
public void onCompleteResponse(ChatResponse completeResponse) {
String finalReply = fullResponse.toString();
String historyItem = "User: " + userQuery + "\nAssistant: " + finalReply; // 构造历史记录格式(根据你项目的实际类型决定,如果是 String 列表,就直接用字符串)
future.complete(Map.of(
SupportState.AGENT_RESPONSE, finalReply, // 同时返回 AGENT_RESPONSE 和 HISTORY,触发 appender 自动追加
SupportState.HISTORY, historyItem
));
}

@Override
public void onError(Throwable error) {
future.completeExceptionally(error);
}
});
return future;
}
}


业务接口类改造

在控制器里,使用 RunnableConfig.builder() 将流接收器注入进去:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
import com.owlias.demo06.state.SupportState;
import jakarta.annotation.Resource;
import org.bsc.async.AsyncGenerator;
import org.bsc.langgraph4j.CompiledGraph;
import org.bsc.langgraph4j.NodeOutput;
import org.bsc.langgraph4j.RunnableConfig;
import org.springframework.http.MediaType;
import org.springframework.scheduling.concurrent.CustomizableThreadFactory;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import reactor.core.publisher.Flux;
import java.util.Map;
import java.util.concurrent.*;

@RestController
public class SupportAgentController {

@Resource
private CompiledGraph<SupportState> compiledGraph;

/**
* 声明一个专属的线程池(通常在 Spring @Configuration 中配置)用于 CompletableFuture.runAsync 异步执行 agent 任务。
* ::::
* CompletableFuture.runAsync 如果不显式传入自定义线程池,它默认使用 JVM 级别的 ForkJoinPool.commonPool()(公共线程池)来执行异步任务。
* 这是一个被 JVM 内所有没有指定线程池的并发任务(如 CompletableFuture.supplyAsync、Java 8+ 的 Stream.parallel() 平行流等)共享的全局线程池。
* ForkJoinPool.commonPool 的默认最大活跃线程数限制为:默认线程数 = CPU 核心数 - 1,如果服务器是 4 核,那么默认只有 3 个线程在跑。
* 所以,绝对不能使用这个默认的线程池,原因在于 “全局共享” 和 “线程饥饿”,这是生产环境接口卡死、雪崩的常见诱因。
* 在生产环境中,CompletableFuture.runAsync 任何异步、I/O 阻塞型任务,都必须显式指定自定义线程池,将其与 CPU 密集型的公共线程池彻底隔离。
*/
ExecutorService agentExecutor = new ThreadPoolExecutor(
10, // 核心线程数(根据并发度决定)
50, // 最大线程数
60L, TimeUnit.SECONDS, // 空闲存活时间
new LinkedBlockingQueue<>(1000), // 有界队列(防止撑爆内存)
new CustomizableThreadFactory("agent-pool-"), // 使用 spring 推荐的 ThreadFactory
new ThreadPoolExecutor.CallerRunsPolicy() // 拒绝策略:忙不过来时让 Controller 线程自己跑,起到限流保护作用
);


@GetMapping(value = "/api/agent/ask", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<String> askAgent(@RequestParam String query) {
return Flux.<String>create(sink -> {
// 1. 默认在公共的 ForkJoinPool 异步启动图,坚决不阻塞 Netty 主线程
CompletableFuture<Void> futureTask = CompletableFuture.runAsync(() -> {
try {
Map<String, Object> inputs = Map.of(SupportState.USER_QUERY, query);

// stream(inputs, runnableConfig) 会启动图。由于内部 Node 的 sseSink.next() 会在 LLM 吐字时立刻触发,前端能够第一时间零延迟拿到打字机效果。
AsyncGenerator.Cancellable<NodeOutput<SupportState>> graphStream = compiledGraph.stream(inputs,
RunnableConfig.builder()
.putMetadata("sseSink", sink) // 通过 RunnableConfig 将 sink 透传,方便后续使用 👈🏻
.build());

// 这里只负责消费工作流节点的变动事件
graphStream.forEach(nodeOutput -> {
String nodeName = nodeOutput.node();
SupportState stateSnapshot = nodeOutput.state();
String eventReport = "\n[系统事件] 节点 [%s] 执行完毕,isSafe:%s\n"
.formatted(nodeName, stateSnapshot.isSafe()); // 只在节点开始或结束时,推送一行轻量级事件(格式化成 SSE 标准)
sink.next(eventReport);
});

// 工作流彻底跑完,关闭 SSE 通道
sink.complete();
} catch (Exception e) {
sink.error(e);
}
}, agentExecutor);

// 2. 监听前端连接断开事件!
sink.onCancel(() -> {
// 前端断开时,立即尝试取消 Future 任务,并给线程发送 Interrupt 中断信号
futureTask.cancel(true); // 立即中断正在运行的后台线程,防止资源泄露!
System.out.println("检测到前端连接断开,已向后台线程发送取消信号。");
});
});
}
}


构件图的改造

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
@Configuration
public class AgentGraphConfig {

@Primary
@Bean("localChatModel")
public ChatModel localChatModel() {
// ...
}

@Bean
public StreamingChatModel localStreamingChatModel() {
// ...
}

@Bean
public CompiledGraph<SupportState> customerSupportGraph(
ChatModel model, StreamingChatModel streamingChatModel) throws Exception {

// 1. 初始化有状态图,绑定 State Schema 和构造器
StateGraph<SupportState> graph = new StateGraph<>(SupportState.SCHEMA, SupportState::new);

// 2. 注入节点 (推荐用 node_async 包装以获取最佳并发性能)
graph.addNode("intent_analyze", AsyncNodeAction.node_async(new IntentAnalyzeNode(model)));
// 👉🏻 graph.addNode(id, asyncNodeActionWithConfig) 本质上是需要一个 asyncNodeActionWithConfig,也就是需要一个 (state, config) -> {... return completableFuture}
graph.addNode("tech_support", (state, config) -> new TechSupportNode(streamingChatModel).execute(state, config));
// 👉🏻 graph.addNode(id, asyncNodeAction) 本质上是需要一个 asyncNodeAction
graph.addNode("sales_support", AsyncNodeAction.node_async(new SalesSupportNode(model)));
graph.addNode("quality_guard", AsyncNodeAction.node_async(new QualityGuardNode()));

// ...
}
}


测试验证

先问一个销售问题,发现得到结果的方式还是和之前相同。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
$ curl -N -G 'http://localhost:8080/api/agent/ask' \
--data-urlencode 'query=请问你们产品打折吗'

data:
data:[系统事件] 节点 [__START__] 执行完毕,isSafe:true
data:

data:
data:[系统事件] 节点 [intent_analyze] 执行完毕,isSafe:true
data:

data:
data:[系统事件] 节点 [sales_support] 执行完毕,isSafe:true # 阻塞等待一次性吐出这个节点执行的结果
data:

data:
data:[系统事件] 节点 [quality_guard] 执行完毕,isSafe:true
data:

data:
data:[系统事件] 节点 [__END__] 执行完毕,isSafe:true
data:

再问一个技术问题,发现 TECH 节点变成了响应式输出:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
$ curl -N -G 'http://localhost:8080/api/agent/ask' \
--data-urlencode 'query=耳机没声音技术上是什么原因'

data:
data:[系统事件] 节点 [__START__] 执行完毕,isSafe:true
data:

data:
data:[系统事件] 节点 [intent_analyze] 执行完毕,isSafe:true
data:

data:您 # 这里开始不再是一次性吐出来 [TECH] 的结果!

data:好

data:!

data:耳机

data:没有

data:声音

data:可能是

data:多种

data:原因

data:造成的
...
...

data:请

data:联系

data:我们

data:更换

data:或

data:维修

data:。

data:
data:[系统事件] 节点 [tech_support] 执行完毕,isSafe:true # [TECH] onCompleteResponse 完成,继续往下走,之后流程和之前相同!
data:

data:
data:[系统事件] 节点 [quality_guard] 执行完毕,isSafe:true
data:

data:
data:[系统事件] 节点 [__END__] 执行完毕,isSafe:true
data:


检查点挂起与人工审核

上述案例的最后可能需要人工审核的环节我们还没有真正实现 。为了在 quality_guard 节点判定 “不合规” 时,流式响应能优雅地停住,并将图的状态挂起等待人工处理,我们需要引入工作流的持久化检查点(Checkpoint Saver)机制。为了方便看到人工审核后,后续节点的执行效果,我们增加一个审核后的节点 response_formatter,它的作用就是对最终安全的内容进行统一封装,确保对客户输出的体验。


更新 Graph

AgentGraphConfig

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
import com.owlias.demo06.node.*;
import com.owlias.demo06.state.SupportState;
import dev.langchain4j.model.chat.ChatModel;
import dev.langchain4j.model.chat.StreamingChatModel;
import dev.langchain4j.model.openai.OpenAiChatModel;
import dev.langchain4j.model.openai.OpenAiStreamingChatModel;
import org.bsc.langgraph4j.CompileConfig;
import org.bsc.langgraph4j.CompiledGraph;
import org.bsc.langgraph4j.GraphRepresentation;
import org.bsc.langgraph4j.StateGraph;
import org.bsc.langgraph4j.action.AsyncEdgeAction;
import org.bsc.langgraph4j.action.AsyncNodeAction;
import org.bsc.langgraph4j.checkpoint.MemorySaver;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Primary;
import java.time.Duration;
import java.util.Map;

@Configuration
public class AgentGraphConfig {

@Primary
@Bean("localChatModel")
public ChatModel localChatModel() {
return OpenAiChatModel.builder() // 👈🏻
.baseUrl("http://localhost:11434/v1")
.apiKey("api_key_xxx")
.modelName("gemma3:1b")
.timeout(Duration.ofSeconds(300))
.maxRetries(3)
.logRequests(true)
.logResponses(true)
.build();
}

@Bean
public StreamingChatModel localStreamingChatModel() {
return OpenAiStreamingChatModel.builder() // 👈🏻
.baseUrl("http://localhost:11434/v1")
.apiKey("api_key_xxx")
.modelName("gemma3:1b")
.timeout(Duration.ofSeconds(300))
.logRequests(true)
.logResponses(true)
.build();
}

@Bean
public CompiledGraph<SupportState> customerSupportGraph(
ChatModel model, StreamingChatModel streamingChatModel) throws Exception {

// 1. 初始化有状态图,绑定 State Schema 和构造器
StateGraph<SupportState> graph = new StateGraph<>(SupportState.SCHEMA, SupportState::new);

// 2. 注入节点 (推荐用 node_async 包装以获取最佳并发性能)
graph.addNode("intent_analyze", AsyncNodeAction.node_async(new IntentAnalyzeNode(model)));
graph.addNode("tech_support", new TechSupportNode(streamingChatModel)::execute);
graph.addNode("sales_support", AsyncNodeAction.node_async(new SalesSupportNode(model)));
graph.addNode("quality_guard", AsyncNodeAction.node_async(new QualityGuardNode()));
// 占位节点:人工审核节点内部通常为空,或者只用来记录“进入审核状态”
graph.addNode("human_review", AsyncNodeAction.node_async(state -> Map.of()));
graph.addNode("response_formatter", AsyncNodeAction.node_async(new ResponseFormatterNode(model)));


// 3. 配置连线:入口点指向意图分析
graph.addEdge(StateGraph.START, "intent_analyze");

// 4. 条件路由分支:
graph.addConditionalEdges("intent_analyze",
// 根据意图路由到技术还是销售
AsyncEdgeAction.edge_async(state -> {
String intent = state.getIntent();
if (intent.contains("TECH")) {
return "route_to_tech";
} else {
return "route_to_sales";
}
}),
// 明确逻辑返回值与节点 ID 的映射关系
Map.of(
"route_to_tech", "tech_support",
"route_to_sales", "sales_support"
)
);

// 5. 业务处理完成后,强制汇总流向质量卫士节点
graph.addEdge("tech_support", "quality_guard");
graph.addEdge("sales_support", "quality_guard");

// 6. 第二个条件分支:如果安全则结束,不安全则走向人工审批/阻断
graph.addConditionalEdges("quality_guard",
AsyncEdgeAction.edge_async(state -> {
if (Boolean.TRUE.equals(state.isSafe())) {
return "safe";
} else {
return "to_review";
}
}),
Map.of(
"safe", "response_formatter", // 👈🏻 安全直接去格式化,而不再走 StateGraph.END!
"to_review", "human_review"
)
);
graph.addEdge("human_review", "response_formatter");
graph.addEdge("response_formatter", StateGraph.END);

// 7. 编译生成可执行的 CompiledGraph
MemorySaver memorySaver = new MemorySaver();
CompileConfig compileConfig = CompileConfig.builder()
.checkpointSaver(memorySaver)
// 👈🏻 声明在这一个节点前需要被挂起/打断,这就好比公路上的红绿灯设在 human_review 节点的之前。
// 当 quality_guard 发现回答不安全时,图引擎就把车开到了 human_review 的路口,看到红灯,刹车挂起。
// 这时候有两件事非常关键:第一,车没过去,human_review 节点自身的代码此时完全没有执行。
// 第二,指针停在门口:此时通过 getState() 查看图的快照,它的 next() 节点明确写着 "human_review"。
// 这就表示:“下一次变绿灯启动时,车子第一步就要开进 human_review”。
.interruptBefore("human_review")
.build();
CompiledGraph<SupportState> compiledGraph = graph.compile(compileConfig);
System.out.println("\n" + compiledGraph.getGraph(GraphRepresentation.Type.MERMAID) + "\n"); ////
return compiledGraph;
}
}

新增的 ResponseFormatterNode:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
/**
* 节点 5:最终响应格式化与礼貌封装节点
* 处于所有分支的终点,对最终安全的内容进行统一封装,确保对客户输出的体验。
*/
public class ResponseFormatterNode implements NodeAction<SupportState> {
private final ChatModel chatModel;

public ResponseFormatterNode(ChatModel chatModel) {
this.chatModel = chatModel;
}

@Override
public Map<String, Object> apply(SupportState state) {
String rawResponse = state.getAgentResponse();

// 如果内容为空,给一个兜底
if (rawResponse == null || rawResponse.trim().isEmpty()) {
return Map.of(SupportState.AGENT_RESPONSE, "抱歉,系统未能生成有效的回复。");
}

// 运用轻量级模型对内容进行润色和结构封装
String prompt = "你是一位专业的客户服务排版与润色助手。请在不改变原意的前提下,"
+ "将以下回答或者后台审核内容润色得更加礼貌、专业,并在文末统一附带一句温馨的企业问候语"
+ "(例如:感谢您的咨询,Owlias 竭诚为您服务!)。\n\n"
+ "待润色的内容为:\n" + rawResponse;

String formattedResponse = chatModel.chat(prompt);

// 更新 AGENT_RESPONSE 为润色后的最终版
return Map.of(SupportState.AGENT_RESPONSE, formattedResponse);
}
}


业务接口类的改造

SupportAgentController 的改造:

  • /api/agent/ask:原提问接口增加一个判断逻辑 “human_review”.equals(nodeName)
  • 暴露两个简单的审核接口 review-test01 和 review-test02(不同的响应方式)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
import com.owlias.demo06.state.SupportState;
import jakarta.annotation.Resource;
import org.bsc.async.AsyncGenerator;
import org.bsc.langgraph4j.*;
import org.bsc.langgraph4j.state.StateSnapshot;
import org.springframework.http.MediaType;
import org.springframework.http.ResponseEntity;
import org.springframework.scheduling.concurrent.CustomizableThreadFactory;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestParam;
import org.springframework.web.bind.annotation.RestController;
import reactor.core.publisher.Flux;
import java.util.HashMap;
import java.util.Map;
import java.util.Optional;
import java.util.concurrent.*;

@RestController
public class SupportAgentController {

@Resource
private CompiledGraph<SupportState> compiledGraph;

ExecutorService agentExecutor = new ThreadPoolExecutor(
10, // 核心线程数(根据并发度决定)
50, // 最大线程数
60L, TimeUnit.SECONDS, // 空闲存活时间
new LinkedBlockingQueue<>(1000), // 有界队列(防止撑爆内存)
new CustomizableThreadFactory("agent-pool-"),
new ThreadPoolExecutor.CallerRunsPolicy() // 拒绝策略:忙不过来时让 Controller 线程自己跑,起到限流保护作用
);


@GetMapping(value = "/api/agent/ask", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<String> askAgent(@RequestParam String query, @RequestParam String userChatId) { // 绑定的用户的会话ID(线程ID)
return Flux.<String>create(sink -> {
CompletableFuture<Void> futureTask = CompletableFuture.runAsync(() -> {
try {
Map<String, Object> inputs = Map.of(SupportState.USER_QUERY, query);
AsyncGenerator.Cancellable<NodeOutput<SupportState>> graphStream = compiledGraph.stream(inputs,
RunnableConfig.builder()
.threadId(userChatId) // 绑定 sseSink 并传入 threadId,用于图的持久化定位 👈🏻
.putMetadata("sseSink", sink) // 通过 RunnableConfig 将 sink 透传,方便后续使用 👈🏻
.build());

graphStream.forEach(nodeOutput -> {
if (Thread.currentThread().isInterrupted()) return;

String nodeName = nodeOutput.node();
SupportState stateSnapshot = nodeOutput.state();

// ⭐️ 改造点:如果走到了人类审核节点,说明 quality_guard 判定不合规
if ("human_review".equals(nodeName)) {
sink.next("\n[警告] 节点 [quality_guard] 判定内容不合规!工作流已挂起,等待人工介入重构...\n");
sink.next("{\"status\": \"SUSPENDED\", \"bad_content\": \"" + stateSnapshot.getAgentResponse() + "\"}");
return;
}

String eventReport = "\n[系统事件] 节点 [%s] 执行完毕,isSafe:%s\n"
.formatted(nodeName, stateSnapshot.isSafe());
sink.next(eventReport);
});

// 工作流彻底跑完,关闭 SSE 通道
sink.complete();
} catch (Exception e) {
sink.error(e);
}
}, agentExecutor);

sink.onCancel(() -> {
futureTask.cancel(true); // 立即中断正在运行的后台线程,防止资源泄露!
System.out.println("检测到前端连接断开,已向后台线程发送取消信号。");
});
});
}

/**
* 极简单次恢复接口:人工干预后静默结束流转
* @param userChatId 绑定的用户的会话ID(使用用户ID也不合适)
* @param approvedContent 审核意见
*/
@GetMapping("/api/agent/review-test01")
public ResponseEntity<String> reviewAndResume1(@RequestParam String userChatId,
@RequestParam String approvedContent) {
try {
RunnableConfig config = RunnableConfig.builder().threadId(userChatId).build();

// 1. 防御式:检测断点是否存在
Optional<StateSnapshot<SupportState>> snapshotOpt = compiledGraph.stateOf(config);
if (snapshotOpt.isEmpty()) {
return ResponseEntity.badRequest().body("未找到该会话对应的可恢复挂起断点。");
}

StateSnapshot<SupportState> snapshot = snapshotOpt.get();
String nextNode = snapshot.next();
if (StateGraph.END.equals(nextNode) || !"human_review".equals(nextNode)) {
return ResponseEntity.badRequest().body("当前会话不处于可审核状态(当前节点: %s)".formatted(nextNode));
}

// 2. 注入人工审核内容
Map<String, Object> currentStateMap = new HashMap<>(snapshot.state().data());
currentStateMap.put(SupportState.AGENT_RESPONSE, approvedContent);
currentStateMap.put(SupportState.IS_SAFE, true); // 人工纠正安全标记

// 3. 提交状态修改,使图复活并指明修改点在 human_review
compiledGraph.updateState(config, currentStateMap, "human_review");

// 4. 恢复断点。使用 GraphInput.resume() 明确告知引擎这是恢复断点,而不是新入参!
compiledGraph.invoke(GraphInput.resume(), config);

return ResponseEntity.ok("人工介入成功,工作流已从断点恢复并顺利结束。");
} catch (Exception e) {
return ResponseEntity.status(500).body("恢复工作流失败: " + e.getMessage());
}
}


/**
* 流式恢复接口:人工审批后,前端不仅能恢复,还能看到后续 formatting 等节点运行的进度
*/
@GetMapping(value = "/api/agent/review-test02", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<String> reviewAndResume2(@RequestParam String userChatId,
@RequestParam String approvedContent) {
return Flux.<String>create(sink -> {

// 核心防御1:主动做断点存在性检查,不把未知的风险带入异步线程池
RunnableConfig baseConfig = RunnableConfig.builder().threadId(userChatId).build();
Optional<StateSnapshot<SupportState>> snapshotOpt = compiledGraph.stateOf(baseConfig);
if (snapshotOpt.isEmpty()) {
// 优雅降级处理:通知前端,人工审核链接已失效
sink.next("\n[警告] 未找到当前会话的挂起断点!\n");
sink.next("[原因] 该会话可能已经流转完毕、审核过了,或是服务器重启导致内存中的断点丢失。\n");
sink.next("[建议] 请让用户重新发起提问,或前往控制台检查会话状态。");
sink.complete();
return;
}

// 核心防御2:如果工作流已经跑完了(nextNode 为空或 __END__),或者已经不卡在 human_review 了
StateSnapshot<SupportState> snapshot = snapshotOpt.get();
String nextNode = snapshot.next();
if (StateGraph.END.equals(nextNode) || !"human_review".equals(nextNode)) {
sink.next("\n[拒绝] 审批失败!该申请之前已经处理完成,或处于 [%s] 状态,不可重复流转。\n".formatted(nextNode));
sink.complete(); // 直接拦截返回,不触发后续的 updateState 和 stream
return;
}

// 将耗时的大模型图流转逻辑,放进自定义的线程池中异步跑
// 统一配置:将此次审核会话专用的 sink 挂载到最新的 activeConfig 中
RunnableConfig activeConfig = RunnableConfig.builder(baseConfig)
.putMetadata("sseSink", sink)
.build();
CompletableFuture<Void> futureTask = CompletableFuture.runAsync(() -> {
try {
// 1. 获取当前被挂起那一刻的最原始快照
Map<String, Object> currentStateMap = new HashMap<>(snapshot.state().data());

// 2. 注入人工审核的结果并修正标记
currentStateMap.put(SupportState.AGENT_RESPONSE, approvedContent);
currentStateMap.put(SupportState.IS_SAFE, true);

// 3. 👉🏻 告诉图:这是由 human_review 节点更新的(图更新状态与恢复流转两个步骤,完美共享同一个前端 SSE 物理连接 sink)
compiledGraph.updateState(activeConfig, currentStateMap, "human_review");
sink.next("[系统提示] 人工审核已通过,正在为您恢复工作流...\n");

// 4. 👉🏻 核心:使用 stream + GraphInput.resume() 恢复流转,并捕获后续节点
AsyncGenerator.Cancellable<NodeOutput<SupportState>> resumeStream = compiledGraph.stream(
GraphInput.resume(),
activeConfig
);

resumeStream.forEach(nodeOutput -> {
if (Thread.currentThread().isInterrupted()) return;

String nodeName = nodeOutput.node();
SupportState stateSnapshot = nodeOutput.state();

// 打印后续流转到的节点事件(比如 response_formatter)
String eventReport = "\n[恢复流转] 节点 [%s] 执行完毕!\n".formatted(nodeName);
sink.next(eventReport);

// 如果是最终润色封装节点,直接把最终答案推给前端
if ("response_formatter".equals(nodeName)) {
sink.next("response_formatter 处理后的响应:" + stateSnapshot.getAgentResponse());
}
});

// 后续节点执行完毕,正常关闭通道
sink.complete();
} catch (Exception e) {
sink.next("\n[错误] 恢复工作流失败: " + e.getMessage() + "\n");
sink.error(e);
}
}, agentExecutor); // 👈 生产必须指定定义的业务线程池!

// 监听前端连接断开,及时中断后台线程,防止线程池泄漏
sink.onCancel(() -> {
futureTask.cancel(true);
System.out.println("检测到审核页面连接断开,已向后台恢复线程发送取消信号。");
});
});
}
}


测试验证

提一个 [TECH] 类的敏感问题,走到 quality_guard 被卡住走审批了(isSafe:false)。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
$ curl -N -G 'http://localhost:8080/api/agent/ask' \
--data-urlencode 'userChatId=1' \
--data-urlencode 'query=公司门禁密码是什么'

data:
data:[系统事件] 节点 [__START__] 执行完毕,isSafe:true
data:

data:
data:[系统事件] 节点 [intent_analyze] 执行完毕,isSafe:true
data:

data:您

data:好

data:,
...

data:公司

data:门

data:禁

data:密码

data:通常

data:是
...

data:
data:[系统事件] 节点 [tech_support] 执行完毕,isSafe:true
data:

data:
data:[系统事件] 节点 [quality_guard] 执行完毕,isSafe:false
data:

当调用审批接口,审批通过时:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
$ curl -N -G 'http://localhost:8080/api/agent/review-test02' \
--data-urlencode 'userChatId=1' \
--data-urlencode 'approvedContent=公司门禁密码是个人账号,如果忘记,重新申请吧。'

data:[系统提示] 人工审核已通过,正在为您恢复工作流...
data:

data:
data:[恢复流转] 节点 [response_formatter] 执行完毕!
data:

data:response_formatter 处理后的响应:好的,下面是对您提供的回答/后台审核内容的润色:
data:
data:**润色后的版本:**
data:
data:“我们理解您可能遇到忘记密码的情况,这确实可能带来一些不便。公司门禁密码是个人账号,如果您忘记了密码,建议您重新申请。 我们的支持团队会及时为您提供帮助,请随时联系我们,我们将竭诚为您服务。”
data:
data:**温馨企业问候语:**
data:
data:感谢您的咨询,Owlias 竭诚为您服务!

data:
data:[恢复流转] 节点 [__END__] 执行完毕!
data:

当再次调用同一个审批接口:

1
2
3
4
5
6
$ curl -N -G 'http://localhost:8080/api/agent/review-test02' \
--data-urlencode 'userChatId=1' \
--data-urlencode 'approvedContent=公司门禁密码是个人账号,如果忘记,重新申请吧。'
data:
data:[拒绝] 审批失败!该申请之前已经处理完成,或处于 [__END__] 状态,不可重复流转。
data:

其次,我们再提一个关于 [SALES] 的正常的问题,没有触发敏感信息,不需要走审批:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
$ curl -N -G 'http://localhost:8080/api/agent/review-test02' \
--data-urlencode 'userChatId=2' \
--data-urlencode 'approvedContent=给我介绍一下这款手机的功能'

data:
data:[系统事件] 节点 [__START__] 执行完毕,isSafe:true
data:

data:
data:[系统事件] 节点 [intent_analyze] 执行完毕,isSafe:true
data:

data:
data:[系统事件] 节点 [sales_support] 执行完毕,isSafe:true
data:

data:
data:[系统事件] 节点 [quality_guard] 执行完毕,isSafe:true
data:

data:
data:[系统事件] 节点 [response_formatter] 执行完毕,isSafe:true
data:

data:
data:[系统事件] 节点 [__END__] 执行完毕,isSafe:true
data:

整个断点审批的闭环,本质上是由前端、Controller 接口和图引擎通过三步完成的:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
 【 第一步:用户提问 】
/api/agent/ask


[quality_guard] (不安全!)


👉🏻 拦截在 human_review 前
(输出 SUSPENDED JSON 给前端)


【 第二步:管理员审批修改 】
/api/agent/review-test02

├─► 1. 拦截重复提交 (检查 next 是否为 human_review)
├─► 2. 修改内存快照 (updateState: 强刷数据, 洗白 isSafe)
└─► 3. 挥动绿灯放行 (stream + GraphInput.resume())


【 第三步:图复活向后流转 】
[human_review] (空节点直接划过)


[response_formatter] (拿着洗白后的干净数据进行最终格式化)


[END] (正式完结)


检查点持久化实现

持久化的选型

要将 LangGraph4J 的检查点(Checkpoint)和状态真正持久化到 MySQL 中,我们需要打破原有的 “内存模式(Memory Saver)”,通过实现 LangGraph4J 的 BaseCheckpointStore 接口,将其对接至数据库持久化API中。这样,即便服务器重启,用户的长对话、或者处于挂起(SUSPENDED)状态等待人工审核的图,都能随时从数据库中精准恢复。

LangGraph 的持久化需要存储两类核心数据:

  • checkpoints:图在每个节点执行完后的逻辑快照(包含版本、下步路由等元数据)。
  • writes:图在当前节点执行期间所产生的状态变更(用于在回滚或复杂重试时重放)。

LangGraph 的检查点(Checkpoint)和状态数据有着非常鲜明的写多读少特征:

  • 图每经过一个节点,甚至节点内每一次状态更新,都会产生一次写入。
  • 大部分时候,引擎只需要通过 thread_id 捞出最新的一条快照,很少需要复杂的关联查询(Join)。
  • 随着对话轮数增加,上下文(Chat History)会越来越大,序列化后的体积不可小觑。

所以,结合这些特点,LangGraph 官方生态,首选的检查点持久化实现是 Postgres SQL(PgSaver)。原因是:

  • PG 强大的 JSONB 支持:不仅能存,还能直接对 State 里的某个 JSON 字段建索引。
  • 优秀的二进制(BYTEA)性能:在高并发写入大二进制块时,PG 的 MVCC 机制和写入性能普遍优于 MySQL。
  • TOAST 技术:PG 会自动将超大的 State 数据压缩并存储在行外,避免大字段污染主表的 Buffer Pool。

如果你的技术栈里已经有 Postgres,或者不介意引入它,无脑选 Postgres。它在关系型数据库的严谨性与 NoSQL 的灵活性之间找到了完美的平衡。当然你也可以选择 redis,特点就是快,但内存存储贵啊,而且持久化可靠性不如 PG。这里我们为了演示的方便,使用 mysql + mybatis 实现。


持久化到 mysql

继续引入 checkpoint mysql 持久化的依赖:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
<dependency>
<groupId>org.bsc.langgraph4j</groupId>
<artifactId>langgraph4j-mysql-saver</artifactId> <!--核心-->
</dependency>

<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-jdbc</artifactId>
</dependency>

<dependency>
<groupId>com.mysql</groupId>
<artifactId>mysql-connector-j</artifactId>
</dependency>

application.yml

1
2
3
4
5
6
7
8
9
10
11
12
spring:
datasource: # 使用 spring jdbc 向容器中注入一个 datasource
url: jdbc:mysql://192.168.1.251:3306/colibri_db?useUnicode=true&characterEncoding=utf8&zeroDateTimeBehavior=convertToNull&useSSL=true&serverTimezone=GMT%2B8
username: xxx
password: xxx
driver-class-name: com.mysql.cj.jdbc.Driver
# 优雅配置:定制连接池参数
hikari:
maximum-pool-size: 10
minimum-idle: 5
idle-timeout: 30000
connection-timeout: 20000

AgentGraphConfig 的改造:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
@Configuration
public class AgentGraphConfig {

@Bean(name = "mysqlSaver")
public MysqlSaver mysqlSaver(DataSource dataSource) {
return new MysqlSaver.Builder().dataSource(dataSource).build();
}

@Bean
public CompiledGraph<SupportState> customerSupportGraph(
ChatModel model, StreamingChatModel streamingChatModel, MysqlSaver mysqlSaver) throws Exception {
// ...

// 7. 编译生成可执行的 CompiledGraph
CompileConfig compileConfig = CompileConfig.builder()
.checkpointSaver(mysqlSaver) // 👈 直接塞入官方的 mysqlSaver 实现 checkpoint 持久化
.interruptBefore("human_review")
.build();
CompiledGraph<SupportState> compiledGraph = graph.compile(compileConfig);
System.out.println("\n" + compiledGraph.getGraph(GraphRepresentation.Type.MERMAID) + "\n"); ////
return compiledGraph;
}
}

然后启动程序,我们观察数据库,就可以看到多了两张表:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
+------------------------+             +-------------------------------+
| langraph4j_thread | | langraph4j_checkpoint |
+------------------------+ +-------------------------------+
| PK | thread_id | <--------+ | PK | checkpoint_id |
| | thread_name | | | FK | thread_id |
| | is_released | | | | node_id / next_node_id |
+------------------------+ +--| | state_data (JSON) |
+-------------------------------+

CREATE TABLE `langraph4j_thread` (
`thread_id` varchar(36) COLLATE utf8mb4_general_ci NOT NULL,
`thread_name` varchar(255) COLLATE utf8mb4_general_ci DEFAULT NULL,
`is_released` tinyint(1) NOT NULL DEFAULT '0',
PRIMARY KEY (`thread_id`),
UNIQUE KEY `IDX_LANGRAPH4J_THREAD_NAME_RELEASED` (`thread_name`,`is_released`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_general_ci;

CREATE TABLE `langraph4j_checkpoint` (
`id` bigint unsigned NOT NULL AUTO_INCREMENT,
`checkpoint_id` varchar(36) COLLATE utf8mb4_general_ci NOT NULL,
`thread_id` varchar(36) COLLATE utf8mb4_general_ci NOT NULL,
`node_id` varchar(255) COLLATE utf8mb4_general_ci DEFAULT NULL,
`next_node_id` varchar(255) COLLATE utf8mb4_general_ci DEFAULT NULL,
`state_data` json NOT NULL, # JSON 类型,需要 mysql 8 的支持
`saved_at` timestamp(6) NULL DEFAULT CURRENT_TIMESTAMP(6),
PRIMARY KEY (`checkpoint_id`),
UNIQUE KEY `id` (`id`),
KEY `LANGRAPH4J_FK_THREAD` (`thread_id`),
CONSTRAINT `LANGRAPH4J_FK_THREAD` FOREIGN KEY (`thread_id`) REFERENCES `langraph4j_thread` (`thread_id`) ON DELETE CASCADE
) ENGINE=InnoDB AUTO_INCREMENT=7 DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_general_ci;

####
# langraph4j_thread (会话主表)
# thread_id:图执行的唯一会话标识。每一个独立的对话(比如一个用户的新聊天窗口)都会分配一个 thread_id。
# is_released:用于标记当前会话是否已经流转结束或处于释放状态。配合 UNIQUE KEY 可以非常优雅地控制同一个命名任务在未释放前不能重复开启。

####
# langraph4j_checkpoint (检查点流水表)
# checkpoint_id:每一次节点执行完毕,都会产生一个全局唯一的 checkpoint_id(通常是 UUID)。
# node_id 与 next_node_id:核心路由标记。记录当前快照是从哪一个节点(Node)执行完荡过来的,以及下一步(Next)即将会流向哪里。这正是断点续传(Resume)的关键。
# state_data (JSON):这是核心。图的所有状态(State Channel 里的变量、对话历史记录等)在这里被完整封包。
# 外键级联 (ON DELETE CASCADE):当你在业务上删除某条会话(langraph4j_thread)时,MySQL 会自动在底层把该会话产生的所有历史检查点全部干净地冲掉,避免产生垃圾数据。


持久化原理剖析

状态的快照演进(Append-only 逻辑):当图在执行时,它不是在同一条记录上反复进行 UPDATE,而是采用 Append-only(只追加不修改) 的时间线模式。

  • 用户输入发问,图启动,插入 langraph4j_thread。
  • 进入 Node_A 执行,结束后,引擎将当前最新的 State 序列化为 JSON,并携带 node_id=’Node_A’ 插入一条新的 langraph4j_checkpoint。
  • 进入 Node_B 执行,结束后,再次插入一条新的 langraph4j_checkpoint 记录。
  • 这种保留历史快照的设计允许 Agent 具备 “时间旅行” 的能力。你可以命令大模型回滚到第 3 轮对话的任意一个历史节点重新执行,而数据库里清晰地保留了过去每一个节点的快照切片。

当我们的图配置了 .interruptBefore(“human_review”) 时,整个运行过程的实现原理如下:

1
2
3
[ 执行节点 A ] ──> [ 触发挂起 ] ──> [ 自动序列化 State 并入库 ] ──> [ 线程销毁, 内存释放 ]

[ 重新唤醒流程 ] <── [ 从数据库捞出最新一条 Checkpoint ] <── [ 管理员审批通过 ]
  • 落库休眠:当程序流转到需要人工介入的节点前,图会自动停下。此时,最后一步的执行结果、上下文以及 next_node_id=’human_review’ 会被打包存入 langraph4j_checkpoint 表。随后该线程直接结束,完全不占用任何系统服务器内存。

  • 状态唤醒与恢复:当管理员点击“审批通过”触发回调时,你传入了对应的 thread_id。

  • 官方的 MySQLSaver 会在底层执行一条 SQL:

    1
    2
    3
    SELECT * FROM langraph4j_checkpoint 
    WHERE thread_id = ?
    ORDER BY saved_at DESC LIMIT 1;
  • 引擎拿到这条最新的快照,将 state_data 反序列化还原回 Java 的 State 内存对象,并看到 next_node_id 是 human_review,于是顺理成章地推动图继续往下走。

这种设计利用了 MySQL 外键维系了长会话生命周期,利用 JSON 类型 兼容了图状态 Schema-less(无固定模式、多变)的特性,利用 saved_at DESC 实现了极其廉价的 “取最新快照” 索引寻址。相比于我们自己去手写 MyBatis 还要去构思复杂的日志追踪表,官方的这两张表已经用最少、最优雅的代价,把图的时间旅行、断点续传和死后复活三大核心能力在关系型数据库上完全坐实了。