找回密码
立即注册
搜索
热搜: Java Python Linux Go
发回帖 发新帖

4197

积分

0

好友

541

主题
发表于 8 小时前 | 查看: 3| 回复: 0

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 这类外部中间件会更合适。

九、总结

  1. SharedMemoryRwLock<HashMap>,线程安全的 KV 存储
  2. Runtime:Graph 的宿主,可选择性挂载 SharedMemory
  3. Workspace:多窗口管理,提供 default 与 named Runtime
  4. Graph 间通过 SharedMemory 传递数据,延迟可以控制在 1ms 以内
  5. 后续可以考虑扩展:持久化到 SQLite / Redis、WebSocket 实时推送、分布式共享内存等

云栈社区 可以找到更多关于 Rust 并发、工作流引擎的实战讨论与资源。




上一篇:SCADA系统详解:数据采集监控、功能架构与PLC/DCS区别
下一篇:树莓派线材避坑指南:电源、HDMI、网线、USB与相机排线选购全解析
您需要登录后才可以回帖 登录 | 立即注册

手机版|小黑屋|网站地图|云栈社区 ( 苏ICP备2022046150号-2 )

GMT+8, 2026-8-2 09:35 , Processed in 1.017436 second(s), 41 queries , Gzip On.

Powered by Discuz! X3.5

© 2025-2026 云栈社区.

快速回复 返回顶部 返回列表