# dag-frame **Repository Path**: alibaba/dag-frame ## Basic Information - **Project Name**: dag-frame - **Description**: No description available - **Primary Language**: Unknown - **License**: Apache-2.0 - **Default Branch**: main - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2026-06-04 - **Last Updated**: 2026-10-01 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # DagFrame 一个高性能的有向无环图(DAG)执行框架,专为复杂的计算任务编排和执行而设计。 ## 🚀 特性 - **高性能**: 基于线程池的异步执行模型,支持高并发场景 - **灵活的图构建**: 支持多种 DAG 拓扑结构(线性链、扇出扇入、二叉树等) - **类型安全**: 基于 C++17 的强类型系统,编译时检查 - **可扩展**: 支持自定义节点和操作(Op) - **组件化**: 支持实验标志、变量、日志、追踪等组件 - **图对象池**: 支持图对象池,减少重复创建的开销 ## 📋 目录 - [架构设计](#架构设计) - [快速开始](#快速开始) - [性能测试](#性能测试) - [使用示例](#使用示例) - [API 文档](#api-文档) - [构建指南](#构建指南) ## 🏗️ 架构设计 ### 核心组件 ``` ┌─────────────────────────────────────────────────────────────┐ │ DagFrame │ │ ┌───────────────────────────────────────────────────────┐ │ │ │ Graph Object Pool │ │ │ │ ┌─────────┐ ┌─────────┐ ┌─────────┐ │ │ │ │ │ Graph 0 │ │ Graph 1 │ │ Graph N │ │ │ │ │ └─────────┘ └─────────┘ └─────────┘ │ │ │ └───────────────────────────────────────────────────────┘ │ │ ↓ │ │ ┌───────────────────────────────────────────────────────┐ │ │ │ Thread Pool │ │ │ │ ┌─────┐ ┌─────┐ ┌─────┐ ┌─────┐ ┌─────┐ │ │ │ │ │ T1 │ │ T2 │ │ T3 │ │ T4 │ │ TN │ │ │ │ │ └─────┘ └─────┘ └─────┘ └─────┘ └─────┘ │ │ │ └───────────────────────────────────────────────────────┘ │ └─────────────────────────────────────────────────────────────┘ ``` ### 关键概念 - **Node**: 计算节点,包含执行逻辑 - **Graph**: DAG 图,包含节点和边 - **Edge**: 节点之间的连接,定义数据流向 - **DagFrame**: 框架主类,管理图对象池和线程池 - **NodeContext**: 节点执行上下文,管理输入输出 ## 🚀 快速开始 ### 环境要求 - C++17 或更高版本 - CMake 3.10 或更高版本 - 依赖库: - gflags - glog - protobuf - TBB (Intel Threading Building Blocks) - gtest ### 构建步骤 ```bash # 克隆仓库 git clone https://github.com/alibaba/dag-frame.git cd dag-frame # 创建构建目录 mkdir build && cd build # 配置和构建 cmake .. make -j$(nproc 2>/dev/null || sysctl -n hw.ncpu 2>/dev/null || echo 4) # 运行测试 ctest --output-on-failure ``` ### 简单示例 ```cpp #include "frame/dag_frame.h" // 定义节点 class MyNode : public dag_frame::Node { public: explicit MyNode(const dag_frame::NodeConstruction& option) : dag_frame::Node(option) {} void DoWork(dag_frame::NodeContext* context) const override { // 你的业务逻辑 std::cout << "Node executed: " << Name() << std::endl; } }; // 注册节点 REGISTER_NODE(MyNode, my_node) .Input("input", "int", false) .Output("output", "int"); // 定义图构造器 class MyGraph : public dag_frame::GraphConstructor { public: void Build(dag_frame::GraphOption* option) const override { option->AddStartNode(0, "my_node"); } }; // 使用框架 int main() { dag_frame::DagFrameOption option; option.InstallGraphConstructor(); option.graph_num_ = 4; option.thread_num_ = 4; dag_frame::DagFrame frame; frame.Build(option); // 执行图 frame.Execute(0); return 0; } ``` ## 📊 性能测试 ### 测试环境 - **CPU**: Intel Xeon (多核) - **操作系统**: Linux (kernel 3.10.0) - **编译器**: GCC 9.2.1 - **C++ 标准**: C++17 - **优化级别**: -O3 ### 测试配置 - **线程数**: 1 - **图对象池大小**: 4 (DagFrame) / 不适用 (Taskflow) - **预热轮数**: 100 - **测试轮数**: 1000 - **节点类型**: 空操作 (No-op) 节点,用于测量纯框架调度开销 ### 测试场景 三种 DAG 拓扑结构: 1. **最简拓扑 (Minimal)**: 2 个节点 (start -> end),测量框架基础开销 2. **线性链 (Linear Chain)**: 10 个节点串行执行 (n0 -> n1 -> ... -> n8 -> end),测量节点间传递开销 3. **扇出扇入 (Fan-out/Fan-in)**: 12 个节点 (start -> n0..n9 -> end),测量并行分发和汇聚开销 ### Benchmark 结果 #### DagFrame vs Taskflow 3.11.0 对比 | 拓扑 | 节点数 | DagFrame 平均延迟 | Taskflow 平均延迟 | DagFrame P50 | Taskflow P50 | DagFrame P99 | Taskflow P99 | |------|--------|------------------|------------------|-------------|-------------|-------------|-------------| | Binary Tree | 2 | 5.67 μs | 8.73 μs | 4.88 μs | 8.73 μs | 14.83 μs | 12.29 μs | | Linear Chain (len=10) | 10 | 12.60 μs | 11.89 μs | 12.00 μs | 11.45 μs | 19.30 μs | 16.29 μs | | Fan-out/Fan-in (width=10) | 12 | 24.17 μs | 15.44 μs | 23.28 μs | 15.39 μs | 38.92 μs | 20.12 μs | #### 吞吐量对比 | 拓扑 | DagFrame (ops/sec) | Taskflow (ops/sec) | |------|-------------------|-------------------| | Binary Tree | ~176,426 | ~114,546 | | Linear Chain (len=10) | ~79,347 | ~84,078 | | Fan-out/Fan-in (width=10) | ~41,376 | ~64,762 | #### DagFrame 详细结果 ##### 1. 二叉树 DAG (Binary Tree) ``` === BinaryTree Benchmark === Rounds: 1000 Performance Results: Average Latency: 5.67 us Min Latency: 3.96 us Max Latency: 30.13 us P50 Latency: 4.88 us P95 Latency: 12.55 us P99 Latency: 14.83 us Throughput: ~176426 ops/sec Total Time: 5.67 ms ``` ##### 2. 线性链 DAG (Linear Chain) ``` === LinearChain (length=10) Benchmark === Rounds: 1000 Performance Results: Average Latency: 12.60 us Min Latency: 11.05 us Max Latency: 299.88 us P50 Latency: 12.00 us P95 Latency: 12.88 us P99 Latency: 19.30 us Throughput: ~79347 ops/sec Total Time: 12.60 ms ``` ##### 3. 扇出扇入 DAG (Fan-out/Fan-in) ``` === FanOutFanIn (width=10) Benchmark === Rounds: 1000 Performance Results: Average Latency: 24.17 us Min Latency: 21.27 us Max Latency: 47.04 us P50 Latency: 23.28 us P95 Latency: 30.65 us P99 Latency: 38.92 us Throughput: ~41376 ops/sec Total Time: 24.17 ms ``` #### Taskflow 3.11.0 详细结果 ##### 1. 二叉树 DAG (Binary Tree) ``` === BinaryTree Benchmark === Rounds: 1000 Performance Results: Average Latency: 8.73 us Min Latency: 2.68 us Max Latency: 17.31 us P50 Latency: 8.73 us P95 Latency: 10.36 us P99 Latency: 12.29 us Throughput: ~114546 ops/sec Total Time: 8.73 ms ``` ##### 2. 线性链 DAG (Linear Chain) ``` === LinearChain (length=10) Benchmark === Rounds: 1000 Performance Results: Average Latency: 11.89 us Min Latency: 9.92 us Max Latency: 20.35 us P50 Latency: 11.45 us P95 Latency: 14.54 us P99 Latency: 16.29 us Throughput: ~84078 ops/sec Total Time: 11.89 ms ``` ##### 3. 扇出扇入 DAG (Fan-out/Fan-in) ``` === FanOutFanIn (width=10) Benchmark === Rounds: 1000 Performance Results: Average Latency: 15.44 us Min Latency: 10.88 us Max Latency: 23.84 us P50 Latency: 15.39 us P95 Latency: 18.11 us P99 Latency: 20.12 us Throughput: ~64762 ops/sec Total Time: 15.44 ms ``` ### 性能分析 **对比总结**: 1. **最简拓扑 (Binary Tree, 2 节点)**: DagFrame (5.67 μs) 优于 Taskflow (8.73 μs),DagFrame 在最小 DAG 上的基础调度开销更低 2. **线性链 (10 节点)**: 两者接近,Taskflow (11.89 μs) 略优于 DagFrame (12.60 μs) 3. **扇出扇入 (12 节点)**: Taskflow (15.44 μs) 明显优于 DagFrame (24.17 μs),Taskflow 的 work-stealing 调度器在并行分支场景下效率更高 4. **尾延迟**: Taskflow 的 P99 延迟整体更稳定 (max 20.12 μs vs 38.92 μs),DagFrame 偶有离群值 (LinearChain max 299.88 μs) **DagFrame 优化方向**: 1. **扇出扇入调度优化**: 当前 DagFrame 在并行分支汇聚时开销较大,可考虑引入 work-stealing 策略 2. **尾延迟控制**: 减少离群值的出现频率 3. **图对象池**: DagFrame 独有的图对象池设计可在高频重复执行场景下发挥优势 4. **异步执行**: 使用 `AsyncNode` 处理 I/O 密集型任务 ### 运行性能测试 ```bash cd build # 运行二叉树 benchmark(支持自定义参数) ./binary_tree_benchmark --benchmark_rounds=1000 --warmup_rounds=100 --thread_num=1 --graph_num=4 # 运行综合性能测试(包含三种拓扑) ./comprehensive_benchmark ``` ## 💡 使用示例 ### 示例 1: 简单的两节点图 ```cpp #include "frame/dag_frame.h" class StartNode : public dag_frame::Node { public: explicit StartNode(const dag_frame::NodeConstruction& option) : dag_frame::Node(option) {} void DoWork(dag_frame::NodeContext* context) const override { int value = 42; context->SetOutput(0, &value); } }; class EndNode : public dag_frame::Node { public: explicit EndNode(const dag_frame::NodeConstruction& option) : dag_frame::Node(option) {} void DoWork(dag_frame::NodeContext* context) const override { const int& value = context->Input(0); std::cout << "Received value: " << value << std::endl; } }; REGISTER_NODE(StartNode, start) .Output("out", "int"); REGISTER_NODE(EndNode, end) .Input("in", "int", false); class SimpleGraph : public dag_frame::GraphConstructor { public: void Build(dag_frame::GraphOption* option) const override { option->AddStartNode(0, "start"); option->AddEdge(0, "start", 0, "end", 0); } }; ``` ### 示例 2: 带条件分支的图 ```cpp class ConditionNode : public dag_frame::Node { public: explicit ConditionNode(const dag_frame::NodeConstruction& option) : dag_frame::Node(option) {} void DoWork(dag_frame::NodeContext* context) const override { int value = context->Input(0); bool condition = value > 50; context->SetConditionValue("condition", condition); } }; REGISTER_NODE(ConditionNode, condition) .Input("in", "int", true); // 在 GraphConstructor 中使用条件 void Build(dag_frame::GraphOption* option) const override { option->AddStartNode(0, "start"); option->AddEdge(0, "start", 0, "condition", 0); // 根据条件选择不同的路径 option->AddEdge(0, "condition", 0, "branch_a", 0); option->AddEdge(0, "condition", 0, "branch_b", 0); } ``` ### 示例 3: 异步节点 ```cpp class AsyncProcessNode : public dag_frame::AsyncNode { public: explicit AsyncProcessNode(const dag_frame::NodeConstruction& option) : dag_frame::AsyncNode(option) {} void AsyncDoWork(dag_frame::NodeContext* context, Callback done) const override { // 在新线程中执行 I/O 操作 std::thread([context, done = std::move(done)]() { // 模拟异步操作 std::this_thread::sleep_for(std::chrono::milliseconds(10)); int result = 123; context->SetOutput(0, &result); // 完成后调用 done done(); }).detach(); } }; REGISTER_NODE(AsyncProcessNode, async_process) .Input("in", "int", false) .Output("out", "int"); ``` ## 📚 API 文档 ### 核心类 #### DagFrame 框架主类,管理图对象池和执行。 ```cpp class DagFrame { public: void Build(const DagFrameOption& option); void Destroy(); template int32_t Execute(size_t graph_index, Targs&&... args); }; ``` #### Node 节点基类,用户需要继承此类实现自定义节点。 ```cpp class Node { protected: virtual void DoWork(NodeContext* context) const; virtual void CustomPreWork(NodeContext* context) const; virtual void CustomPostWork(NodeContext* context) const; }; ``` #### GraphConstructor 图构造器基类,用于定义 DAG 拓扑结构。 ```cpp class GraphConstructor { public: virtual void Build(GraphOption* option) const = 0; }; ``` ### 宏定义 #### REGISTER_NODE 注册节点到框架。 ```cpp REGISTER_NODE(ClassName, node_name) .Input("input_name", "type", required) .Output("output_name", "type") .Attr("attr_name", default_value); ``` ### 工具类 #### BenchmarkStats 性能统计工具。 ```cpp struct BenchmarkStats { double avg_us = 0; // 平均延迟(微秒) double min_us = 0; // 最小延迟 double max_us = 0; // 最大延迟 double p50_us = 0; // P50 延迟 double p95_us = 0; // P95 延迟 double p99_us = 0; // P99 延迟 double throughput = 0; // 吞吐量(ops/sec) double total_us = 0; // 总耗时(微秒) int rounds = 0; // 测试轮数 }; // 计算统计数据 BenchmarkStats ComputeStats(std::vector& latencies); // 打印结果 void PrintResults(const std::string& name, const BenchmarkStats& stats); ``` ## 🔧 构建指南 ### CMake 选项 ```bash # 启用调试模式 cmake -DCMAKE_BUILD_TYPE=Debug .. # 启用优化模式 cmake -DCMAKE_BUILD_TYPE=Release .. # 指定安装路径 cmake -DCMAKE_INSTALL_PREFIX=/usr/local .. ``` ### 运行测试 项目使用 Google Test 框架进行单元测试。在构建完成后,可以使用以下命令运行测试: #### 运行所有测试 ```bash cd build ctest --output-on-failure ``` 这将运行所有测试用例,并在测试失败时显示详细的错误信息。 #### 运行特定测试套件 ```bash # 运行核心库测试 ./dag_core_test # 运行框架测试 ./dag_frame_test # 运行示例测试 ./example_binary_test ./example_dag_frame_test # 运行线程池测试 ./thread_pool_test ``` #### 运行特定测试用例 ```bash # 运行特定的测试用例 ./dag_core_test --gtest_filter=GraphTest.Graph0Test ./dag_core_test --gtest_filter=GraphTest.Graph1Test # 运行特定测试套件的所有测试 ./dag_core_test --gtest_filter=NodeTest.* ./dag_core_test --gtest_filter=GraphTest.* # 运行多个测试用例(使用冒号分隔) ./dag_core_test --gtest_filter=GraphTest.Graph0Test:GraphTest.Graph1Test ``` #### 运行性能测试 ```bash cd build # 运行二叉树 benchmark(支持自定义参数) ./binary_tree_benchmark --benchmark_rounds=1000 --warmup_rounds=100 --thread_num=1 --graph_num=4 # 运行综合性能测试(包含三种拓扑) ./comprehensive_benchmark ``` #### 测试覆盖率 查看测试覆盖率(需要安装 gcov 和 lcov): ```bash # 重新构建时启用覆盖率 mkdir -p build_coverage && cd build_coverage cmake -DCMAKE_BUILD_TYPE=Debug -DENABLE_COVERAGE=ON .. make -j$(nproc 2>/dev/null || sysctl -n hw.ncpu 2>/dev/null || echo 4) # 运行测试 ctest # 生成覆盖率报告 lcov --capture --directory . --output-file coverage.info lcov --remove coverage.info '/usr/*' '*/third_party/*' '*/test/*' --output-file coverage.info genhtml coverage.info --output-directory coverage_report ``` #### 依赖安装 #### macOS ```bash brew install gflags glog protobuf tbb googletest ``` #### Linux (Ubuntu/Debian) ```bash sudo apt-get install -y libgflags-dev libgoogle-glog-dev libprotobuf-dev protobuf-compiler libtbb-dev libgtest-dev ``` #### Linux (CentOS/RHEL) ```bash sudo yum install -y gflags-devel glog-devel protobuf-devel protobuf-compiler tbb-devel gtest-devel ``` ## 🤝 贡献指南 欢迎贡献!请遵循以下步骤: 1. Fork 本仓库 2. 创建特性分支 (`git checkout -b feature/AmazingFeature`) 3. 提交更改 (`git commit -m 'Add some AmazingFeature'`) 4. 推送到分支 (`git push origin feature/AmazingFeature`) 5. 开启 Pull Request ## 📄 许可证 本项目采用 Apache License 2.0 许可证 - 详见 [LICENSE](LICENSE) 文件 ## 📞 联系方式 - 项目主页: [https://github.com/alibaba/dag-frame](https://github.com/alibaba/dag-frame) - 问题反馈: [https://github.com/alibaba/dag-frame/issues](https://github.com/alibaba/dag-frame/issues) ## 🙏 致谢 感谢所有贡献者的支持! --- **注意**: 本项目正在积极开发中,API 可能会发生变化。建议在生产环境使用前进行充分测试。