微信扫码
添加专属顾问
Spring AI Alibaba Graph 1.0.0.4版本重磅升级,流式输出全面拥抱响应式编程,大幅提升开发效率! 核心内容: 1. 传统迭代器模式与响应式流模式的架构差异分析 2. 基于Project Reactor的响应式流实现方案详解 3. 新版流式输出机制的性能优化与集成优势
20250925日,随着Spring AI Alibaba Graph从1.0.0.3升级至1.0.0.4,其中的Graph流式输出有了很大的改进,相关的example已更新,欢迎大家随时跟进,PR地址如下:https://github.com/spring-ai-alibaba/examples/pull/364
随着 spring-ai-alibaba-graph 模块的广泛应用,社区中出现了许多关于其流式处理实现机制的疑问和使用文档需求。其中最突出的问题是:graph 模块的流式实现如何与当前主流的响应式流(Reactive Streams)框架进行集成。
当前 graph 模块采用传统的迭代器模式实现流式输出,这种实现方式与主流的响应式编程范式存在较大差异。在与 Project Reactor、RxJava 等响应式流框架集成时,开发者需要进行大量额外的适配工作,增加了技术复杂性和维护成本。
为解决这一问题并提升框架的现代化程度,团队对 graph 内核中的流式输出机制进行了重构,采用基于 Project Reactor 的响应式流实现,以更好地与现代响应式生态系统集成,降低开发者的使用门槛并提高整体性能表现。
传统迭代器模式在 Spring AI Alibaba 项目的 AsyncGenerator<E> 接口实现流式输出。其核心设计思路是将异步数据流抽象为一个可迭代的对象,消费者通过调用 next() 方法逐个获取数据元素。
该模式遵循以下原则:
Data<E> 类封装异步操作的各种状态(正常数据、完成状态、错误状态)传统迭代器模式更符合命令式编程思维,而响应式流模式则体现了声明式编程的理念。
传统迭代器模式主要依赖以下 Java 并发机制:
thenCompose、thenApply 等方法构建异步操作链// AsyncGenerator接口中的toCompletableFuture方法体现了链式调用思想
default CompletableFuture<Object> toCompletableFuture() {
final Data<E> next = next();
if (next.isDone()) {
return completedFuture(next.resultValue);
}
return next.data.thenCompose(v -> toCompletableFuture());
}AsyncGenerator<E> 接口定义了异步生成器的核心契约:
public interface AsyncGenerator<E> extends Iterable<E> {
Data<E> next(); // 获取下一个异步数据元素
default CompletableFuture<Object> toCompletableFuture() { ... } // 转换为CompletableFuture
default Stream<E> stream() { ... } // 转换为Stream
default Iterator<E> iterator() { ... } // 获取迭代器
}主要实现类包括:
Data<E> 类是异步数据元素的核心封装,设计上考虑了多种状态:
class Data<E> {
final CompletableFuture<E> data; // 异步数据
final Embed<E> embed; // 嵌入式生成器
final Object resultValue; // 结果值
publicbooleanisDone() { // 完成状态判断
return data == null && embed == null;
}
publicbooleanisError() { // 错误状态判断
return data != null && data.isCompletedExceptionally();
}
}设计考量包括:
class WithResult<E> implementsAsyncGenerator<E>, HasResultValue {
protectedfinal AsyncGenerator<E> delegate;
private Object resultValue;
@Override
publicfinal Data<E> next() {
final Data<E> result = delegate.next();
if (result.isDone()) {
resultValue = result.resultValue;
}
return result;
}
}class WithEmbed<E> implements AsyncGenerator<E>, HasResultValue {
protected final Deque<Embed<E>> generatorsStack = new ArrayDeque<>(2);
private final Deque<Data<E>> returnValueStack = new ArrayDeque<>(2);
@Override
public Data<E> next() {
// 处理嵌套生成器栈
// 实现生成器组合逻辑
}
}传统迭代器模式的背压处理主要通过以下方式实现:
next() 方法调用会阻塞局限性:
AsyncGenerator 继承 Iterable 接口WithResult 和 WithEmbed 类next() 方法定义算法框架Spring AI Alibaba Graph 模块的流式处理设计采用了响应式编程范式,基于 Project Reactor 框架实现。其核心理念是:
OverAllState 管理流式处理过程中的状态变化Flux 作为流式处理的核心组件具有以下优势:
在 ParallelNode 中,流合并策略采用以下设计:
Flux.zip 操作符合并多个 Flux 流OverAllState.updateState 方法维护状态一致性为保持向后兼容性,项目提供了 AsyncGenerator 与 Flux 的双向转换:
AsyncGenerator.fromFlux:将 Flux 转换为 AsyncGeneratorFlowGenerator.fromPublisher:将 Publisher 转换为 AsyncGenerator在 ParallelNode 中,多个 Flux 流的合并通过以下步骤实现:
// 检查是否有Flux类型的输出
booleanhasFlux= results.stream()
.flatMap(map -> map.values().stream())
.anyMatch(value -> value instanceof Flux);
if (hasFlux) {
// 收集所有Flux流
List<Flux<Object>> fluxList = newArrayList<>();
// ... 处理非Flux输出 ...
// 合并Flux流
if (!fluxList.isEmpty()) {
Flux<Object> mergedFlux = Flux.zip(fluxList, newFunction<Object[], Object>() {
@Override
public Object apply(Object[] objects) {
returnnull; // 简化的合并逻辑
}
});
mergedState.put("__merged_stream__", mergedFlux);
}
}StreamingOutput 类封装了流式输出的数据:
public classStreamingOutputextendsNodeOutput {
privatefinal String chunk;
privatefinal ChatResponse chatResponse;
publicStreamingOutput(ChatResponse chatResponse, String node, OverAllState state) {
super(node, state);
this.chatResponse = chatResponse;
this.chunk = null;
}
publicStreamingOutput(String chunk, String node, OverAllState state) {
super(node, state);
this.chunk = chunk;
this.chatResponse = null;
}
}StreamingChatGenerator 构建流式聊天生成器:
public AsyncGenerator<? extendsNodeOutput> buildInternal(Flux<ChatResponse> flux,
Function<ChatResponse, StreamingOutput> outputMapper) {
varresult=newAtomicReference<ChatResponse>(null);
Consumer<ChatResponse> mergeMessage = (response) -> {
result.updateAndGet(lastResponse -> {
// 合并消息逻辑
// ...
});
};
varprocessedFlux= flux
.filter(response -> response.getResult() != null && response.getResult().getOutput() != null)
.doOnNext(mergeMessage)
.map(next -> newStreamingOutput(next.getResult().getOutput().getText(), startingNode, startingState));
return FlowGenerator.fromPublisher(FlowAdapters.toFlowPublisher(processedFlux),
() -> mapResult.apply(result.get()));
}OverAllState 通过以下方式管理流式处理状态:
public static Map<String, Object> updateState(Map<String, Object> state, Map<String, Object> partialState,
Map<String, KeyStrategy> keyStrategies) {
Objects.requireNonNull(state, "state cannot be null");
if (partialState == null || partialState.isEmpty()) {
return state;
}
Map<String, Object> updatedPartialState = updatePartialStateFromSchema(state, partialState, keyStrategies);
return Stream.concat(state.entrySet().stream(), updatedPartialState.entrySet().stream())
.collect(toMapRemovingNulls(Map.Entry::getKey, Map.Entry::getValue, (currentValue, newValue) -> newValue));
}AsyncGeneratorUtils 提供多生成器合并功能:
public static <T> AsyncGenerator<T> createMergedGenerator(List<AsyncGenerator<T>> generators,
Map<String, KeyStrategy> keyStrategyMap) {
returnnewAsyncGenerator<>() {
privatefinalStampedLocklock=newStampedLock();
privateAtomicIntegerpollCounter=newAtomicInteger(0);
private Map<String, Object> mergedResult = newHashMap<>();
privatefinal List<AsyncGenerator<T>> activeGenerators = newCopyOnWriteArrayList<>(generators);
@Override
public Data<T> next() {
// 轮询各个生成器,合并结果
// ...
}
};
}在并行执行中,流式数据通过以下方式传递和聚合:
CompletableFuture.allOf 并行执行多个节点OverAllState.updateState 更新全局状态__merged_stream__ 键用于标识合并后的 Flux 流,在后续处理中可以识别和处理合并的流数据。
项目通过类型检查来处理不同数据类型:
instanceof Flux 检查,收集到 fluxList 中进行合并OverAllState.updateState 方法更新状态onBackpressureBuffer() 处理背压CompletableFuture 实现异步执行graph TD
A[StateGraph] --> B[ParallelNode]
B --> C[AsyncParallelNodeAction]
C --> D[Node Actions]
D --> E[Flux Streams]
E --> F[Flux Merge]
F --> G[__merged_stream__]
G --> H[OverAllState]sequenceDiagram
participant Client
participant ParallelNode
participant NodeAction1
participant NodeAction2
participant FluxMerge
Client->>ParallelNode: Execute
ParallelNode->>NodeAction1: Execute Async
ParallelNode->>NodeAction2: Execute Async
NodeAction1-->>ParallelNode: Return Flux
NodeAction2-->>ParallelNode: Return Flux
ParallelNode->>FluxMerge: Merge Flux Streams
FluxMerge-->>ParallelNode: Merged Flux
ParallelNode->>Client: Return Resultnext() 调用可能涉及线程阻塞响应式流模式在资源利用方面具有明显优势:
53AI,企业落地大模型首选服务商
产品:场景落地咨询+大模型应用平台+行业解决方案
承诺:免费POC验证,效果达标后再合作。零风险落地应用大模型,已交付160+中大型企业
2026-09-08
我把微软的本体学习工具 Ontology-Playground 魔改了
2026-09-08
企业架构文档转本体模型,本体模型如何更好的解释和支撑企业LTC端到端流程
2026-09-06
万字长文:缺算力还是语义?企业Agent的真正瓶颈
2026-09-04
拆解 Claude Code 记忆系统:渐进式加载 + 双向链接构成的 Markdown 知识图谱
2026-09-01
VisActor 全新图可视化开源项目:VGraph
2026-08-31
构建基于WorkBuddy的元技能-实现基于本体模型驱动的领域技能动态创建
2026-08-24
从RAG到GraphRAG:用Neo4j构建知识图谱, 让知识库真正"活"起来
2026-08-18
文本如何变成知识图谱:两条方法路线与三个开源项目
2026-08-03
2026-06-25
2026-06-24
2026-06-24
2026-07-27
2026-08-04
2026-07-01
2026-07-22
2026-07-01
2026-07-31
欢迎您使用【53AI 官方网站】(以下简称“本网站”或“我们”)。本《会员服务协议》(以下简称“本协议”)是您(以下简称“会员”或“用户”)与【深圳市博思协创网络科技有限公司】之间关于注册、登录及使用本网站会员服务所订立的法律协议。
在您注册或登录前,请务必审慎阅读、充分理解各条款内容,特别是免除或限制责任的条款、知识产权条款、争议解决条款等。此类条款将以加粗形式提示您注意。 当您通过微信公众号授权、手机验证码验证或其他方式成功登录本网站时,即视为您已完全理解并同意接受本协议的全部内容。
一、 定义
本网站:指由【深圳市博思协创网络科技有限公司】运营的,域名为【53ai.com】的网站及相关移动端页面。
会员服务:指本网站向注册会员提供的知识库文章查阅、内容检索及其他相关增值服务。
知识库内容:指本网站发布的包括但不限于文字、图表、数据、研究报告、行业分析等数字化内容资源。
二、 账号注册与登录
登录方式:本网站支持以下登录方式,您可根据实际情况选择:
微信公众号授权登录:您同意将您的微信OpenID信息授权给本网站,用于创建或关联会员账号。
手机验证码登录:您需提供真实有效的手机号码,并通过短信验证码完成身份验证与登录/注册。
账号安全:您的账号仅限您本人使用,禁止赠与、借用、租用、转让或售卖。因您保管不善导致的账号被盗、密码泄露等损失,由您自行承担。
实名认证:根据相关法律法规要求,我们可能要求您在特定功能下完成实名认证。如您拒绝提供,可能无法使用部分或全部服务。
未成年人保护:若您未满18周岁,请在法定监护人的陪同下阅读本协议,并在征得监护人同意后使用本服务。
三、 服务内容与规范
知识库查阅权限:会员登录后,有权按照其会员等级对应的权限范围,在线浏览、检索本网站知识库中的相关文章及内容。
服务变更:我们有权根据业务发展需要,调整、变更或终止部分服务内容,并将以网站公告、公众号消息等方式提前通知。
禁止行为:您在使用服务时不得实施以下行为:
利用技术手段批量爬取、下载、转存知识库内容;
将知识库内容用于商业目的或未经授权地向第三方传播;
干扰本网站正常运行或侵犯其他用户合法权益;
发布违法违规信息或从事违反公序良俗的活动。
四、 知识产权声明
权利归属:本网站知识库中的排版设计、软件代码等内容的知识产权均归【公司全称】或原权利人所有,受《中华人民共和国著作权法》等法律保护。
有限许可:本网站授予会员一项非独占、不可转让、不可转授权的普通许可,仅限于个人学习、研究之目的在线查阅知识库内容。
侵权追责:未经书面许可,任何单位或个人不得以任何形式复制、转载、摘编、镜像、汇编或以其他方式使用上述内容。一经发现,我们保留追究其法律责任的权利。
五、 个人信息保护
我们重视对您个人信息的保护。关于我们如何收集、使用、存储和保护您的个人信息,请单独阅读 《隐私政策》。
您通过微信公众号授权或手机号验证所提供的信息,我们将严格按照《个人信息保护法》的规定处理,仅用于身份识别、服务提供及安全验证等必要用途。
您可以随时通过网站设置或联系客服行使查阅、更正、删除个人信息及撤回授权同意的权利。
六、 免责声明
内容准确性:知识库内容仅供参考,不构成专业建议。我们不对其完整性、准确性、时效性作任何明示或暗示的保证,您应自行判断并承担使用风险。
不可抗力:因自然灾害、政策法规变化、网络故障、第三方平台接口异常(如微信接口维护、运营商短信通道故障)等不可抗力导致的服务中断或延迟,我们不承担违约责任。
第三方链接:本网站可能包含指向第三方网站的链接,该等网站的内容和服务不受我们控制,请您自行甄别风险。
七、 违约责任
如您违反本协议约定,我们有权视情节采取警告、限制功能、暂停服务、注销账号等措施,并保留要求赔偿损失的权利。
如因您的违约行为导致我们遭受行政处罚、第三方索赔或商誉损失,您应承担全部赔偿责任(包括但不限于罚款、赔偿金、律师费、公证费等)。
八、 法律适用与争议解决
本协议的订立、执行和解释均适用中华人民共和国大陆地区法律。
因本协议产生的或与本协议有关的任何争议,双方应友好协商解决;协商不成的,任何一方均可向【公司所在地】有管辖权的人民法院提起诉讼。
九、 其他
本协议构成双方就本服务达成的完整协议,取代此前任何口头或书面约定。
本协议任一条款被认定为无效或不可执行的,不影响其他条款的效力。
我们对本协议享有最终解释权,并在法律允许的范围内保留随时修改的权利。修改后的协议一经公布即生效,继续使用服务即视为同意修订内容。