# OSS.PipeLine **Repository Path**: osscore/OSS.PipeLine ## Basic Information - **Project Name**: OSS.PipeLine - **Description**: 流式事件处理,微服务下业务生命周期管理,强化业务的流程管理,建立业务操作边界,打造标准化的业务执行单元,提高代码复用。 - **Primary Language**: Unknown - **License**: GPL-3.0 - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 4 - **Forks**: 3 - **Created**: 2021-06-01 - **Last Updated**: 2026-08-10 ## Categories & Tags **Categories**: process-engine **Tags**: None ## README # OSS.Pipeline 轻量级 .NET 流程引擎与数据管道框架,用于构建异步数据处理管道和业务流程。 ## 项目组成 | 库 | NuGet | 描述 | |---|---|---| | **OSS.DataPipe** | `Install-Package OSS.DataPipe` | 轻量级消息管道库,支持生产-消费模式和基于消息管道的自动重试事件管理 | | **OSS.Pipeline** | `Install-Package OSS.Pipeline` | 轻量级流程引擎,提供 Activity、Branch、Convertor、Enumerator 组件 | ## 架构关系 ``` ┌─────────────────────────────────────────────────────────────┐ │ OSS.Pipeline │ │ (Workflow Engine - Activity, Branch, Convertor, Enumerator)│ └───────────────────────┬─────────────────────────────────────┘ │ depends on ▼ ┌─────────────────────────────────────────────────────────────┐ │ OSS.DataPipe │ │ (Message Pipe - Producer/Consumer with Retry Support) │ └─────────────────────────────────────────────────────────────┘ ``` --- ## OSS.DataPipe 快速开始 ### 基本使用 ```csharp using OSS.DataPipe; // 推荐使用静态全局定义(内部使用默认内存队列) private static readonly IDataProducer producer = DataPipeFactory.CreateProducer(async (data) => { Console.WriteLine($"处理数据: {data}"); return true; // 返回 true 表示成功,false 触发重试 }); await producer.Push("测试数据"); producer.Release(); // 仅自定义 Provider 时需要释放 ``` ### 自定义队列 Provider 实现 `IDataPipeProvider` 接口可集成 RabbitMQ、Kafka 等消息队列: ```csharp internal class CustomDataPipeProvider : IDataPipeProvider { public IDataProducer? CreateProducer(IDataConsumer consumer, DataPipeOption option) { return option?.SourceCode switch { "RabbitMQ" => new RabbitMqProducer(consumer), _ => null // 返回 null 使用默认内存队列 }; } } DataPipeFactory.PipeProvider = new CustomDataPipeProvider(); ``` ### 事件重试 (RetryProcessor) ```csharp using OSS.DataPipe; // 自动重试 3 次 public static readonly RetryProcessor _processor = new(DoSomething, new RetryPolicy(3)); var result = await _processor.Process("业务参数"); // result.processed_flag: Success/Retrying/Failure // result.executed_times: 已执行次数 ``` 详见 [RetryProcessor 详细文档](docs/DataPipe.RetryProcessor.md) --- ## OSS.Pipeline 快速开始 每个业务流程由多个业务事件组成,每个事件当做一个 Pipe。上游 Pipe 的输出作为下游 Pipe 的输入,类型一致即可串联。 ### 组件类型 | 组件 | PipeType | 描述 | |------|----------|------| | **Activity** | `Activity` | 基本执行单元,处理业务逻辑 | | **BranchGateway** | `BranchGateway` | 条件分支,根据条件路由到不同路径(并发执行) | | **Convertor** | `Convertor` | 类型转换器,桥接不同类型 | | **Enumerator** | `Enumerator` | 枚举器,将集合展开为单个元素 | ### 构建简单流程 ```csharp using OSS.Pipeline; var startPipe = new SimpleActivity("step1", async (input) => { Console.WriteLine($"{input} -> 步骤1完成"); return "Y"; }); var endPipe = startPipe .Append("step2", async (resStr) => { Console.WriteLine($"{resStr} -> 步骤2完成"); return resStr == "Y" ? 10 : 0; }) .Append("step3", async (count) => { Console.WriteLine($"{(count == 10 ? "成功" : "失败")} -> 步骤3完成"); return string.Empty; }); var pipeline = PipelineBuilder.Build("简单流程", startPipe, endPipe); var result = await pipeline.Run("测试流程"); ``` ### 分支流程 ```csharp var branchGate = new BranchPipe("审核结果分支"); createActivity.Append(reduceInventoryActivity).Append(auditActivity).Append(branchGate); // 分支1:审核通过 branchGate.Append(result => result.IsSuccess(), new NotifyEmailActivity()); // 分支2:审核失败 - 内联定义 branchGate.Append("退回库存", result => !result.IsSuccess(), async (result) => { await ReturnInventoryAsync(result); return new NotifyResult("已退回库存"); }); var pipeline = PipelineBuilder.Build("订单流程", createActivity, tailActivity); ``` ### 类型转换器 ```csharp var convertor = orderActivity.Append( order => new OrderDto { Id = order.Id, Amount = order.Amount }, "OrderToDto"); ``` ### 配置组件重试 ```csharp var activity = new SimpleActivity("外部API", CallApi) .SetRetryPolicy(new RetryPolicy(3, autoReleaseDataPipe: true)); ``` ### 流程监控 ```csharp public class MyMonitor : IPipeMonitor { public Task Monitor(MonitorDataItem data) { Console.WriteLine($"[{data.Stage}] {data.PipeCode}"); // data.PipeType, data.Input, data.Output, data.ExecutedTimes, data.Exception return Task.CompletedTask; } } pipeline.SetMonitor(new MyMonitor()); ``` ### 自定义 Activity ```csharp public class CreateOrderActivity : BaseActivity { public CreateOrderActivity() : base("CreateOrder") { } protected override async Task> Executing( CreateOrderReq para, CancellationToken ct) { var order = await _orderService.CreateAsync(para); return new EventResult(order); // 成功 // return new EventResult(true, "失败"); // 需要重试 } } ``` --- ## 应用场景 | 场景 | DataPipe | Pipeline | 典型用例 | |------|:--------:|:--------:|---------| | 消息队列处理 | ✓ | | 异步消息消费、跨服务解耦 | | 事件驱动架构 | ✓ | ✓ | 事件发布/订阅、事件溯源 | | 订单处理流程 | | ✓ | 创建→库存→支付→通知 | | 审批工作流 | | ✓ | 提交→审核→分支→归档 | | 数据同步/ETL | ✓ | ✓ | 读取→转换→写入→缓存 | | 任务编排 | | ✓ | 定时任务调度、依赖管理 | --- ## 构建与测试 ```bash dotnet build src/OSS.PipeLine.sln dotnet run --project src/Tests/OSS.Pipeline.ConsoleTests/OSS.Pipeline.ConsoleTests.csproj ``` ## 目录结构 ``` src/ ├── OSS.DataPipe/ # 消息管道库 │ ├── Interface/ # 核心接口 (IDataProducer/Consumer/Provider) │ ├── DefaultQueue/ # 默认内存队列 (InterQueueHub/Producer) │ └── Event/ # 重试系统 (RetryProcessor/IRetryEvent) ├── OSS.Pipeline/ # 流程引擎库 │ ├── Base/ # 基础抽象 (BasePipeMeta/PipeType) │ ├── Pipeline/ # 流程构建器 & 监控 │ └── Pipes/ # 组件实现 │ ├── Activity/ # SimpleActivity/BaseActivity/PassiveActivity │ ├── Branch/ # BranchPipe │ ├── Convertor/ # ConvertorPipe │ └── Enumerator/ # EnumeratorPipe └── Tests/ # 控制台测试项目 ``` ## 详细文档 | 文档 | 描述 | |------|------| | [文档导航](docs/README.md) | 完整文档索引 | | [DataPipe.Overview](docs/DataPipe.Overview.md) | 核心接口、DataPipeFactory、配置、快速开始 | | [DataPipe.RetryProcessor](docs/DataPipe.RetryProcessor.md) | RetryProcessor 重试处理器详解 | | [DataPipe.Advanced](docs/DataPipe.Advanced.md) | 自定义队列、DataPipeWraper、性能优化、最佳实践 | | [Pipeline.Overview](docs/Pipeline.Overview.md) | 概述、组件类型、快速开始、核心接口 | | [Pipeline.Components](docs/Pipeline.Components.md) | Activity、Branch、Convertor、Enumerator 详解 | | [Pipeline.Advanced](docs/Pipeline.Advanced.md) | 流程监控、嵌套流程、错误处理、最佳实践 | | [Scenarios](docs/Scenarios.md) | 应用场景与集成示例 | ## 许可证 MIT License