在 Rust 生态中构建多Agent工作流,数据共享是避不开的问题。AutoAgentDraw 选择 SharedMemory 作为 Graph 之间的数据桥梁,基于 RwLock 保证线程安全,实现低延迟的进程内通信。
一、为什么需要共享内存?
在多Agent工作流跑起来之后,最常见的需求:
- Graph 1 分析完数据 → Graph 2 需要读取这个结果
- 多个 Graph 需要共享全局配置或运行状态
- 前端需要实时获取工作流进度
SharedMemory 就是为此而生。
二、架构
Graph 1 数据分析 |
Graph 2 报告生成 |
SharedMemory RwLock<HashMap> set / get / delete / keys |
Workspace default Runtime named Runtime |
三、SharedMemory 实现
核心就是一个 RwLock<HashMap>,读写分离、线程安全,比传统消息队列更轻量:
// 共享内存 - Graph间通信的桥梁
pub struct SharedMemory {
data: parking_lot::RwLock<HashMap<String, serde_json::Value>>,
}
impl SharedMemory {
pub fn new() -> Self {
Self { data: parking_lot::RwLock::new(HashMap::new()) }
}
// 写入数据
pub fn set(&self, key: &str, value: serde_json::Value) {
self.data.write().insert(key.to_string(), value);
}
// 读取数据
pub fn get(&self, key: &str) -> Option<serde_json::Value> {
self.data.read().get(key).cloned()
}
// 删除数据
pub fn delete(&self, key: &str) {
self.data.write().remove(key);
}
// 列出所有键
pub fn keys(&self) -> Vec<String> {
self.data.read().keys().cloned().collect()
}
}
四、Runtime:Graph 的宿主
Runtime 管理多个 Graph 的生命周期,可以挂载一个 共享内存 实例(可选),这样一来 Graph 之间就有了公共的 KV 存储:
// 运行时 - 管理所有Graph的生命周期
pub struct Runtime {
graphs: HashMap<GraphId, Graph>,
shared_memory: Option<Arc<SharedMemory>>,
}
impl Runtime {
// 普通运行时(无共享内存)
pub fn new() -> Self {
Self { graphs: HashMap::new(), shared_memory: None }
}
// 带共享内存的运行时(推荐)
pub fn with_shared_memory() -> Self {
Self { graphs: HashMap::new(), shared_memory: Some(Arc::new(SharedMemory::new())) }
}
// 创建新Graph
pub fn create_graph(&mut self, name: &str) -> GraphId {
let graph = Graph::new(name);
let id = graph.id().clone();
self.graphs.insert(id.clone(), graph);
id
}
// 获取共享内存引用
pub fn shared_memory(&self) -> Option<&Arc<SharedMemory>> {
self.shared_memory.as_ref()
}
}
五、Workspace:多窗口支持
Workspace 用来管理多个 Runtime,支持创建命名的 Runtime 实例,为多窗口或多会话场景提供隔离的数据空间:
// Workspace - 多窗口工作空间
pub struct Workspace {
runtimes: HashMap<String, Arc<RwLock<Runtime>>>,
default: Arc<RwLock<Runtime>>,
}
impl Workspace {
pub fn new() -> Self {
let default = Arc::new(RwLock::new(Runtime::with_shared_memory()));
Self { runtimes: HashMap::new(), default }
}
// 获取默认运行时
pub async fn runtime(&self) -> Arc<RwLock<Runtime>> {
self.default.clone()
}
// 创建命名运行时(多窗口)
pub async fn create_runtime(&mut self, name: &str) -> Arc<RwLock<Runtime>> {
let runtime = Arc::new(RwLock::new(Runtime::with_shared_memory()));
self.runtimes.insert(name.to_string(), runtime.clone());
runtime
}
// 获取命名运行时
pub async fn get_runtime(&self, name: &str) -> Option<Arc<RwLock<Runtime>>> {
self.runtimes.get(name).cloned()
}
}
六、完整演示
现在用 SharedMemory 将 3 个 Graph 串联起来——一个写分析结果,一个读状态,另一个后续处理,都在同一个 Runtime 内发生:
use agentflow::runtime::{Runtime, Workspace};
#[tokio::main]
async fn main() {
// 创建带共享内存的工作空间
let workspace = Workspace::new();
let mut rt = workspace.runtime().await.write().await;
// 创建3个独立工作流
let wf1 = rt.create_graph("数据分析");
let wf2 = rt.create_graph("报告生成");
let wf3 = rt.create_graph("结果审核");
println!("创建了 {} 个工作流", rt.list_graphs().len());
// 共享内存:Graph间传递数据
let mem = rt.shared_memory().unwrap();
mem.set("analysis_result", serde_json::json!({"accuracy": 0.95}));
mem.set("report_status", serde_json::json!("pending"));
println!("共享内存键: {:?}", mem.keys());
println!("分析结果: {:?}", mem.get("analysis_result"));
// 销毁工作流
rt.destroy_graph(&wf3);
println!("销毁后: {} 个工作流", rt.list_graphs().len());
}
七、cargo run 输出
$ cargo run
创建了 3 个工作流
共享内存键: ["analysis_result", "report_status"]
分析结果: Some(Object({"accuracy": Number(0.95)}))
销毁后: 2 个工作流
八、设计对比
| 方案 |
并发安全 |
延迟 |
适用场景 |
| SharedMemory |
RwLock保证 |
< 1ms |
同进程内Graph间通信 |
| MessageBus |
RwLock保证 |
< 1ms |
Agent间异步消息 |
| Redis/外部 |
服务端保证 |
1-10ms |
跨进程/分布式 |
SharedMemory 在延迟上碾压外部方案,但局限是同一进程内;如果需要跨进程或多服务交互,Redis 这类外部中间件会更合适。
九、总结
- SharedMemory:
RwLock<HashMap>,线程安全的 KV 存储
- Runtime:Graph 的宿主,可选择性挂载 SharedMemory
- Workspace:多窗口管理,提供 default 与 named Runtime
- Graph 间通过 SharedMemory 传递数据,延迟可以控制在 1ms 以内
- 后续可以考虑扩展:持久化到 SQLite / Redis、WebSocket 实时推送、分布式共享内存等
在 云栈社区 可以找到更多关于 Rust 并发、工作流引擎的实战讨论与资源。