Workflow Core 2.1.2 新能力解析:使用 IWorkflowPurger 清理持久化存储中的历史工作流
发布时间:2026/10/12 1:42:39 锦皓数字建站

后端工作流自动化流程编排【免费下载链接】workflow-coreLightweight workflow engine for .NET Standard项目地址https://gitcode.com/gh_mirrors/wo/workflow-core点击查看免费下载导读本文基于 Workflow Core 2.1.2 版本记录ReleaseNotes/2.1.2.md深入讲解该版本新增的IWorkflowPurger服务——一个按工作流状态与完成时间批量清除历史实例的清理机制。读完本文你将掌握它的接口语义、IoC 注册方式、各持久化提供商的底层实现差异以及如何在自己的应用中安全地执行周期性清理。从版本记录说起2.1.2 引入了什么Workflow Core 2.1.2 的版本记录内容非常聚焦只宣布了一件事Adds a feature to purge old workflows from the persistence store.即新增了从持久化存储中清除旧工作流的能力。配套引入了一个全新的服务接口IWorkflowPurger它可以从 IoC 容器中直接注入并暴露如下方法签名Task PurgeWorkflows(WorkflowStatus status, DateTime olderThan)该版本记录同时注明当时该接口的实现仅覆盖 SQL Server、PostgreSQL 和 MongoDB 三种持久化提供商。从当前仓库源码看这一能力此后已被更多提供商接入详见下文各持久化提供商的实现差异。在长期运行的工作流系统中每条工作流实例WorkflowInstance在结束后仍会以完整快照形式驻留在持久化存储中包含执行指针ExecutionPointers、扩展属性等大量关联数据。随着时间推移这些历史数据会持续占用存储并拖慢查询性能因此需要一个安全、可控的清理机制——这正是IWorkflowPurger的定位。IWorkflowPurger 接口一次调用按状态与时间批量清理接口定义位于 src/WorkflowCore/Interface/IWorkflowPurger.csusing System; using System.Threading; using System.Threading.Tasks; using WorkflowCore.Models; namespace WorkflowCore.Interface { public interface IWorkflowPurger { Task PurgeWorkflows(WorkflowStatus status, DateTime olderThan, CancellationToken cancellationToken default); } }两个核心参数的含义如下参数类型语义statusWorkflowStatus要清理的工作流状态只清理与该状态匹配的实例olderThanDateTime时间阈值只清理完成时间CompleteTime早于该时间点的实例cancellationTokenCancellationToken可选取消令牌默认default用于中断长时间清理操作WorkflowStatus枚举定义于 src/WorkflowCore/Models/WorkflowInstance.cs共四个取值public enum WorkflowStatus { Runnable 0, // 可运行排队等待执行 Suspended 1, // 挂起等待外部事件或延迟 Complete 2, // 已完成 Terminated 3, // 已终止 }实例的CompleteTime属性DateTime?由执行引擎在实例进入完成或终止状态时写入——在 WorkflowExecutor.cs 中可以看到workflow.Status WorkflowStatus.Complete;与workflow.CompleteTime _datetimeProvider.UtcNow;成对出现。也就是说只有真正跑完或被终止的实例才带有结束时间戳这也是清理判据依赖CompleteTime而非CreateTime的原因。通过 IoC 容器注入与使用默认注册内存提供商的兜底实现调用AddWorkflow()时框架会默认注册一个基于内存存储的IWorkflowPurger实现。注册逻辑位于 src/WorkflowCore/ServiceCollectionExtensions.csservices.AddSingletonISingletonMemoryProvider, MemoryPersistenceProvider(); if (!services.Any(x x.ServiceType typeof(IWorkflowPurger))) { services.AddSingletonIWorkflowPurger(sp (IWorkflowPurger)sp.GetServiceISingletonMemoryProvider()); }即只要没有其他持久化提供商显式注册过IWorkflowPurgerMemoryPersistenceProvidersrc/WorkflowCore/Services/DefaultProviders/MemoryPersistenceProvider.cs就会同时充当 purger因为该类直接实现了IWorkflowPurger接口public class MemoryPersistenceProvider : ISingletonMemoryProvider, IWorkflowPurger。使用 SQL Server / PostgreSQL 等提供商时自动替换当你通过扩展方法注册 SQL Server 或 PostgreSQL 等持久化提供商时对应的WorkflowOptions扩展会自动注册真正的持久化 purger 实现例如 SqlServer/ServiceCollectionExtensions.csoptions.Services.AddTransientIWorkflowPurger(sp new WorkflowPurger(new SqlContextFactory(connectionString, initAction)));因此业务代码无需关心底层存储是什么只需面向接口注入public class WorkflowCleanupService { private readonly IWorkflowPurger _purger; public WorkflowCleanupService(IWorkflowPurger purger) { _purger purger; } public async Task CleanupAsync() { // 清理 30 天前已完成的工作流 await _purger.PurgeWorkflows(WorkflowStatus.Complete, DateTime.UtcNow.AddDays(-30)); // 清理 90 天前已终止的工作流 await _purger.PurgeWorkflows(WorkflowStatus.Terminated, DateTime.UtcNow.AddDays(-90)); } }说明IWorkflowPurger本身并不包含调度能力框架也并未内置周期性清理任务。在生产环境中常见的实践是借助IHostedService、BackgroundService或定时作业如 Hangfire / Quartz按固定周期调用上述方法如果使用 SQL Server 队列提供商也可以参考 QueueProviders.SqlServer 提供的命令队列能力自行编排。这部分属于应用层编排属于对本文能力的合理使用方式而非框架内置功能。各持久化提供商的实现差异PurgeWorkflows的语义在不同存储中保持一致状态匹配 完成时间早于阈值但底层删除策略各不相同。下面按当前仓库源码逐一说明。SQL Server / PostgreSQL / MySQL / SQLite / Oracle基于 Entity Framework这五类提供商共用同一套基于 EF Core 的 purger 实现src/providers/WorkflowCore.Persistence.EntityFramework/Services/WorkflowPurger.cs。public async Task PurgeWorkflows(WorkflowStatus status, DateTime olderThan, CancellationToken cancellationToken default) { var olderThanUtc olderThan.ToUniversalTime(); using (var db ConstructDbContext()) { var workflows await db.SetPersistedWorkflow().Where(x x.Status status x.CompleteTime olderThanUtc).ToListAsync(cancellationToken); foreach (var wf in workflows) { foreach (var pointer in wf.ExecutionPointers) { foreach (var extAttr in pointer.ExtensionAttributes) { db.Remove(extAttr); } db.Remove(pointer); } db.Remove(wf); } await db.SaveChangesAsync(cancellationToken); } }实现要点先将olderThan统一转换为 UTCToUniversalTime()再与存储中的CompleteTime比较避免时区偏差导致误删采用先查后删的策略按Status status CompleteTime olderThanUtc查出符合条件的实例由于工作流实例与执行指针、扩展属性存在一对多关联删除实例前会先逐层移除ExtensionAttributes和ExecutionPointers再删除实例本身最后统一SaveChangesAsync提交事务。注册路径可参见 PostgreSQL/ServiceCollectionExtensions.cs、MySQL/ServiceCollectionExtensions.cs、Sqlite/ServiceCollectionExtensions.cs 以及 Oracle 的对应扩展。MongoDBMongoDB 的实现位于 src/providers/WorkflowCore.Persistence.MongoDB/Services/WorkflowPurger.cs直接使用驱动级批量删除var olderThanUtc olderThan.ToUniversalTime(); await WorkflowInstances.DeleteManyAsync(x x.Status status x.CompleteTime olderThanUtc, cancellationToken);由于 MongoDB 文档天然支持嵌套结构一次DeleteManyAsync即可删除整个实例文档无需手动级联清理子对象。注册时由 MongoDB/ServiceCollectionExtensions.cs 完成支持两种方式直接传mongoUrldatabaseName或传入FuncIServiceProvider, IMongoDatabase自定义数据库实例。RavenDBRavenDB 的实现位于 src/providers/WorkflowCore.Persistence.RavenDB/Services/WorkflowPurger.cs使用服务端 RQL 查询删除var utcTime olderThan.ToUniversalTime(); var queryToDelete new IndexQuery { Query $FROM {nameof(WorkflowInstance)} where status {status} and CompleteTime {olderThan} }; return _database.Operations.SendAsync(new DeleteByQueryOperation(queryToDelete, new QueryOperationOptions { AllowStale false }), token: cancellationToken);通过DeleteByQueryOperation在服务端执行删除AllowStale false保证按索引快照的一致性删除。注册位于 RavenDB/ServiceCollectionExtensions.cs。Azure Cosmos DBCosmos DB 的实现位于 src/providers/WorkflowCore.Providers.Azure/Services/WorkflowPurger.cs采用 LINQ to Feed 迭代 逐条删除的方式using (FeedIteratorPersistedWorkflow feedIterator _workflowContainer.Value.GetItemLinqQueryablePersistedWorkflow() .Where(x x.Status status x.CompleteTime olderThanUtc) .ToFeedIterator()) { while (feedIterator.HasMoreResults) { foreach (var item in await feedIterator.ReadNextAsync(cancellationToken)) { await _workflowContainer.Value.DeleteItemAsyncPersistedWorkflow(item.id, new PartitionKey(item.id), cancellationToken: cancellationToken); } } }符合 Cosmos 的分页读取模型FeedIterator分批拉取再逐个删除注册为单例见 Azure/ServiceCollectionExtensions.cs。内存提供商默认MemoryPersistenceProvider的 purge 实现有两个值得注意的约束MemoryPersistenceProvider.csif (status ! WorkflowStatus.Complete status ! WorkflowStatus.Terminated) { throw new ArgumentOutOfRangeException(nameof(status), status, Only complete or terminated workflows can be purged.); } lock (_instances) { _instances.RemoveAll(x x.Status status x.CompleteTime.HasValue x.CompleteTime.Value olderThanUtc); }只允许Complete或Terminated状态传入Runnable或Suspended会直接抛出ArgumentOutOfRangeException从语义上杜绝误删仍在执行中的实例要求CompleteTime有值CompleteTime为 null即从未真正结束的实例即使状态匹配也不会被清除。语义与边界什么会被清理什么不会综合各实现源码可以归纳出PurgeWorkflows(status, olderThan)的完整行为边界双条件过滤实例必须同时满足Status status且CompleteTime olderThanUTC两个条件缺一不可时间统一按 UTC 比较所有实现都会先对入参调用ToUniversalTime()因此调用方传入本地时间或 UTC 时间均可比较结果以 UTC 为准只清终态实例从内存实现的限制可见该接口的设计意图是清理已结束完成/终止的实例Runnable与Suspended状态即使时间久远也不在清理范围内级联清理关联数据关系型数据库中实例的执行指针与扩展属性会随实例一起被移除避免残留孤儿数据。单元测试佐证仓库中的单元测试直接验证了只清理终态、不影响活跃实例的核心语义见 test/WorkflowCore.UnitTests/Services/MemoryPersistenceProviderFixture.csPurgeWorkflows_should_remove_terminated_instances构造一个CompleteTime在两天前、状态为Terminated的实例和一个Runnable的活跃实例执行PurgeWorkflows(WorkflowStatus.Terminated, DateTime.UtcNow.AddDays(-1))后断言只剩活跃实例存活PurgeWorkflows_should_remove_completed_instances同样的模式验证Complete状态实例被清理而一天前创建的活跃实例不受影响。这两个用例同时也演示了最贴近实战的调用姿势清理阈值如AddDays(-1)应晚于被清理实例的CompleteTime、早于当前时间从而保证只清除真正过期、不再需要的实例。小结IWorkflowPurger是 Workflow Core 2.1.2 引入的一个小而实用的能力它以单一接口统一了按状态 按完成时间清理历史工作流的操作通过 IoC 容器与持久化提供商自动绑定应用层无需感知底层存储差异。从当前仓库可以看到该能力已从最初的 SQL Server / PostgreSQL / MongoDB 扩展到了基于 Entity Framework 的全部关系型提供商以及 RavenDB、Azure Cosmos DB内存提供商则提供了带保护性约束的兜底实现。对于需要长期运行、存储会持续增长的工作流系统而言将 purge 接入定期清理任务是控制存储成本与维持查询性能的推荐做法。赞分享后端工作流自动化流程编排【免费下载链接】workflow-coreLightweight workflow engine for .NET Standard项目地址https://gitcode.com/gh_mirrors/wo/workflow-core点击查看免费下载相关推荐DeepSeek Harness 持久化工作流运行记录让 Chat 中的 workflow 历史跨刷新、跨会话完整可追溯DeepSeek Harness 持久化工作流运行记录让 Chat 中的 workflow 历史跨刷新、跨会话完整可追溯 本篇技术指南以仓库内已实现特性笔记人工智能AI AgentAgent 框架DeepSeekPath of Building中文版从新手到高手的流放之路角色构建指南Path of Building中文版从新手到高手的流放之路角色构建指南 你是否在《流放之路》中苦苦挣扎于复杂的角色构建面对英文版的Path of Buil游戏开发GraphiQL查询历史持久化存储与智能管理的实现方案GraphiQL查询历史持久化存储与智能管理的实现方案 引言为什么查询历史管理如此重要 在GraphQL开发过程中开发者经常需要反复执行相似的查询来测试开发工具后端上一篇TypeScript类型推断FE-Interview中的泛型约束题下一篇Flutter中的定时任务使用WorkManager实现后台周期性任务创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
锦
锦皓数字建站
深耕本土企业品牌数字化升级,专注原创端正雅致商务官网,从视觉设计到稳定运维全程保驾护航。