后端工作流自动化流程编排【免费下载链接】workflow-coreLightweight workflow engine for .NET Standard项目地址https://gitcode.com/gh_mirrors/wo/workflow-core点击查看免费下载WorkflowCore 默认以单节点模式运行其内置的SingleNodeQueueProvider与SingleNodeLockProvider全部基于进程内内存实现。本文基于官方文档 docs/multi-node-clusters.md完整梳理将该工作流引擎横向扩展为多节点集群所必需的外部队列机制与分布式锁管理器并给出各 Provider 的注册方式、源码级原理与关键注意事项帮助你在 Azure、AWS、Redis、RabbitMQ、SQL Server 等基础设施上构建高可用的工作流处理集群。适用前提本文所有 Provider 均以 NuGet 包形式提供且集群模式下必须同时配置外部持久化存储内存持久化无法跨节点共享。注册代码中的services.AddWorkflow(...)均位于Microsoft.Extensions.DependencyInjection命名空间下由 WorkflowCore.ServiceCollectionExtensions 提供。一、为什么单节点配置无法支撑集群WorkflowCore 的宿主服务WorkflowHost在启动时会依次启动队列 Provider、锁 Provider、事件中枢与后台任务见 WorkflowHost.StartAsync。其默认配置来自 WorkflowOptionsQueueFactory new FuncIServiceProvider, IQueueProvider(sp new SingleNodeQueueProvider()); LockFactory new FuncIServiceProvider, IDistributedLockProvider(sp new SingleNodeLockProvider()); PersistenceFactory new FuncIServiceProvider, IPersistenceProvider(sp new TransientMemoryPersistenceProvider(...));这组默认值意味着队列、锁、持久化全部停留在单个进程内部跨节点完全不可见因此无法支撑多节点集群。要组成集群必须分别替换队列与锁两个维度持久化同理。1. 单节点队列进程内 BlockingCollectionSingleNodeQueueProvidersrc/WorkflowCore/Services/DefaultProviders/SingleNodeQueueProvider.cs内部使用DictionaryQueueType, BlockingCollectionstring维护三类队列枚举定义见 IQueueProvider.csQueueType.Workflow0待执行的工作流实例QueueType.Event1待分发的事件QueueType.Index2待写入索引的工作流。IsDequeueBlocking为true出队会阻塞等待新消息由于队列只存在于当前进程内存中其他节点无法感知或消费这些消息。2. 单节点锁进程内 HashSetSingleNodeLockProvidersrc/WorkflowCore/Services/DefaultProviders/SingleNodeLockProvider.cs用一个加锁的HashSetstring模拟互斥。AcquireLock在集合中不存在该 Id 时才将其加入并返回true否则返回falseReleaseLock则将其移除。这套机制只在同一进程内有效多个进程同时调用时没有任何全局互斥语义。3. 锁在消费链路中的关键作用锁并非装饰品。在 WorkflowConsumer.ProcessItem 中每个节点从队列取出工作流 Id 后第一步就是await _lockProvider.AcquireLock(itemId, cancellationToken)获取失败 → 该工作流可能正被其他节点处理当前节点直接返回获取成功 → 加载实例、执行、持久化、重新入队并在finally中ReleaseLock(itemId)。若所有节点共享同一个外部锁如 Redis RedLock、Azure Blob Lease、DynamoDB 条件写同一工作流实例在同一时刻只会被一个节点执行——这正是集群不重复处理、不产生数据竞争的根本保证。测试基类 DistributedLockProviderTests 针对所有分布式锁实现统一验证了可获取锁、重复获取失败、释放后可再获取三个契约行为。4. 队列消费与并发上限后台消费者 QueueConsumer 循环DequeueWork并在每个节点本地控制最大并发数Workflow 队列的并发上限由WorkflowOptions.MaxConcurrentWorkflows决定默认Math.Max(Environment.ProcessorCount, 4)。在集群中各节点独立消费共享队列横向扩容即增加消费者数量但同一时刻单条工作流只会有一个节点持锁执行。二、集群的两大核心组件按官方文档的划分组建集群需要同时配置两类 Provider组件职责默认实现单节点集群可选实现队列 ProviderIQueueProvider分发待处理的工作流 / 事件 / 索引任务SingleNodeQueueProviderAzure Storage Queues、Redis、RabbitMQ、AWS SQS、SQL Server Service Broker分布式锁管理器IDistributedLockProvider保证同一工作流实例同时只被一个节点处理SingleNodeLockProviderAzure Blob Storage Leases、Redis RedLock、AWS DynamoDB两个接口均位于 src/WorkflowCore/InterfaceIQueueProvider.csQueueWork入队、DequeueWork出队、IsDequeueBlocking标识阻塞/非阻塞语义、Start/Stop生命周期IDistributedLockProvider.csAcquireLock(Id, cancellationToken)返回是否获取成功、ReleaseLock(Id)释放、Start/Stop生命周期。经验要点锁 Id 通常就是工作流实例 Id如wf:{id}、事件键evt:{id}等资源标识。从 WorkflowConsumer 可看到同一实例在执行—入队—再次出队—再执行的全过程中反复获取/释放锁锁粒度直接决定集群的并行度。三、队列 Provider 详解与注册方式1. RedisWorkflowCore.Providers.Redis安装dotnet add package WorkflowCore.Providers.Redis注册示例连接字符串localhost:6379app-name为队列键前缀channel-name为事件总线频道名services.AddWorkflow(cfg { cfg.UseRedisPersistence(localhost:6379, app-name); cfg.UseRedisLocking(localhost:6379); cfg.UseRedisQueues(localhost:6379, app-name); cfg.UseRedisEventHub(localhost:6379, channel-name); });源码证据ServiceCollectionExtensions.csUseRedisQueues(connectionString, prefix)→ 注册RedisQueueProviderUseRedisLocking(connectionString, prefix null)→ 注册RedisLockProviderUseRedisPersistence(connectionString, prefix, deleteComplete false)→ 注册RedisPersistenceProviderUseRedisEventHub(connectionString, channel)→ 注册RedisLifeCycleEventHub。实现原理队列基于 Redis List 实现RedisQueueProvider.cs队列名形如{prefix}-workflows/{prefix}-events/{prefix}-index入队使用ListInsertBeforeAsync优先插入出队使用ListLeftPopAsyncIsDequeueBlocking false空队列立即返回null由消费者按Options.IdleTime轮询默认 100ms见 WorkflowOptions.cs锁基于RedLock 算法RedisLockProvider.cs_lockTimeout默认 1 分钟锁资源名形如{prefix}:{key}连接断开时会在再次获取锁前自动重连EnsureConnected。2. RabbitMQWorkflowCore.QueueProviders.RabbitMQ仅提供队列能力需另行搭配分布式锁管理器文档原文明确 along with a distributed lock manager。安装dotnet add package WorkflowCore.QueueProviders.RabbitMQ注册READMEservices.AddWorkflow(x x.UseRabbitMQ(new ConnectionFactory() { HostName localhost }));源码证据ServiceCollectionExtensions.cs提供三个重载UseRabbitMQ(IConnectionFactory connectionFactory)单连接工厂UseRabbitMQ(IConnectionFactory connectionFactory, IEnumerablestring hostnames)多主机地址UseRabbitMQ(RabbitMqConnectionFactory rabbitMqConnectionFactory)自定义连接工厂委托RabbitMqConnectionFactory为delegate TaskIConnection(IServiceProvider sp, string clientProvidedName, CancellationToken)。内部还会注册IRabbitMqQueueNameProvider默认实现DefaultRabbitMqQueueNameProvider用于决定集群各队列的命名。3. Azure Storage QueuesWorkflowCore.Providers.Azure安装dotnet add package WorkflowCore.Providers.Azure注册READMEUseAzureSynchronization同时完成队列 分布式锁两项配置services.AddWorkflow(options { options.UseAzureSynchronization(azure storage connection string); options.UseAzureServiceBusEventHub(service bus connection string, topic name, subscription name); options.UseCosmosDbPersistence(connection string); });源码证据ServiceCollectionExtensions.cspublic static WorkflowOptions UseAzureSynchronization(this WorkflowOptions options, string connectionString) { options.UseQueueProvider(sp new AzureStorageQueueProvider(connectionString, sp.GetServiceILoggerFactory())); options.UseDistributedLockManager(sp new AzureLockManager(connectionString, sp.GetServiceILoggerFactory())); return options; }另有基于UriTokenCredential如DefaultAzureCredential托管标识的重载适合无连接字符串的认证场景持久化可选用 Cosmos DBUseCosmosDbPersistence或更经济的 Azure Table StorageUseAzureTableStoragePersistence(connectionString, tableNamePrefix WorkflowCore)也支持传入TableServiceClient或Uri TokenCredential。事件分发可使用UseAzureServiceBusEventHub基于 Service Bus 的LifeCycleEventHub。4. AWS Simple Queue ServiceWorkflowCore.Providers.AWS安装dotnet add package WorkflowCore.Providers.AWS注册README持久化 队列 分布式锁三项一起配置services.AddWorkflow(cfg { cfg.UseAwsDynamoPersistence(new EnvironmentVariablesAWSCredentials(), new AmazonDynamoDBConfig() { RegionEndpoint RegionEndpoint.USWest2 }, table-prefix); cfg.UseAwsDynamoLocking(new EnvironmentVariablesAWSCredentials(), new AmazonDynamoDBConfig() { RegionEndpoint RegionEndpoint.USWest2 }, workflow-core-locks); cfg.UseAwsSimpleQueueService(new EnvironmentVariablesAWSCredentials(), new AmazonSQSConfig() { RegionEndpoint RegionEndpoint.USWest2 }, queues-prefix); });若 AWS 资源不存在Provider 会自动创建DynamoDB 表与索引默认按吞吐量 1 预置可到 AWS 控制台调整。若已持有预配置客户端可使用WithProvisionedClient系列方法ServiceCollectionExtensions.csvar client new AmazonDynamoDBClient(); var sqsClient new AmazonSQSClient(); services.AddWorkflow(cfg { cfg.UseAwsDynamoPersistenceWithProvisionedClient(client, table-prefix); cfg.UseAwsDynamoLockingWithProvisionedClient(client, workflow-core-locks); cfg.UseAwsSimpleQueueServiceWithProvisionedClient(sqsClient, queues-prefix); });另提供基于 Kinesis 的事件总线UseAwsKinesis(credentials, RegionEndpoint, appName, streamName)其会在 DynamoDB 中创建表跟踪每个 shard 的消费位置。5. SQL Server Service BrokerWorkflowCore.QueueProviders.SqlServer社区贡献感谢 Roberto Paterlini的仅队列Provider基于 SQL Server Service Broker 实现同样需配合分布式锁管理器。安装dotnet add package WorkflowCore.QueueProviders.SqlServer注册READMEservices.AddWorkflow(x x.UseSqlServerBroker(Server.;DatabaseWorkflowCore;Trusted_ConnectionTrue;, true, true));源码证据ServiceCollectionExtensions.csUseSqlServerBroker(connectionString, canCreateDb, canMigrateDb)其中canCreateDb允许创建数据库、canMigrateDb允许迁移队列结构内部注册IQueueConfigProvider、ISqlCommandExecutor、ISqlServerQueueProviderMigrator最终以SqlServerQueueProviderOptions含ConnectionString/CanCreateDb/CanMigrateDb构建SqlServerQueueProvider。四、分布式锁管理器详解与注册方式1. Azure Blob Storage LeasesWorkflowCore.Providers.Azure与 Azure Storage Queues 同包通过UseAzureSynchronization一并注册见上文源码底层使用 Blob Lease 实现锁语义适合与 Azure 全家桶搭配。2. RedisWorkflowCore.Providers.Redis通过UseRedisLocking(connectionString, prefix null)注册底层为 RedLock 算法见 RedisLockProvider.csAcquireLock通过RedLockFactory.CreateLockAsync(GetResource(Id), _lockTimeout)获取锁IsAcquired为真才返回true成功获取的锁被记录进ManagedLocksReleaseLock按资源名匹配并Dispose_lockTimeout固定为 1 分钟锁资源名支持prefix隔离{prefix}:{key}。3. AWS DynamoDBWorkflowCore.Providers.AWS通过UseAwsDynamoLocking(credentials, config, tableName)或UseAwsDynamoLockingWithProvisionedClient(dynamoClient, tableName)注册底层使用 DynamoDB 条件写入实现分布式锁tableName为锁表名示例workflow-core-locks表不存在时自动创建。五、完整集群示例从单节点到三节点以Redis 全家桶为例最小可运行的多节点配置如下每个节点运行相同代码app-name保持一致以保证共享同一组队列键节点间通过connectionString指向同一 Redis 实例public void ConfigureServices(IServiceCollection services) { services.AddWorkflow(cfg { // 1. 持久化所有节点必须共享同一存储不能用内存实现 cfg.UseRedisPersistence(localhost:6379, app-name); // 2. 分布式锁保证同一工作流同时只被一个节点执行 cfg.UseRedisLocking(localhost:6379); // 3. 队列让所有节点从同一组 List 中取任务 cfg.UseRedisQueues(localhost:6379, app-name); // 4. 可选跨节点事件总线 cfg.UseRedisEventHub(localhost:6379, channel-name); }); services.AddHostedServiceWorkflowHostedService(); } public class WorkflowHostedService : BackgroundService { private readonly IWorkflowHost _host; public WorkflowHostedService(IWorkflowHost host) _host host; protected override async Task ExecuteAsync(CancellationToken stoppingToken) { _host.Start(); // 注册工作流定义、启动实例、发布事件等 await Task.Delay(Timeout.Infinite, stoppingToken); } public override async Task StopAsync(CancellationToken cancellationToken) { _host.Stop(); await base.StopAsync(cancellationToken); } }对应其他基础设施的等价组合Azure 集群UseAzureSynchronization(connectionString)UseCosmosDbPersistence(...)或UseAzureTableStoragePersistence(...)可选UseAzureServiceBusEventHub做事件总线AWS 集群UseAwsDynamoPersistence(...)UseAwsDynamoLocking(...)UseAwsSimpleQueueService(...)可选UseAwsKinesis(...)做事件总线混合型例Redis 锁 RabbitMQ 队列 任意外部持久化如 PostgreSQL / SQL Server / MongoDB见 src/providers 下各 Persistence 项目。六、集群配置清单与注意事项三件套缺一不可集群 外部队列 分布式锁 共享持久化。默认的 TransientMemoryPersistenceProvider 的ScheduleCommand直接抛NotImplementedException、SupportsScheduledCommands false且数据仅在单进程内可见任何多节点场景都必须替换持久化实现。锁与队列必须配套RabbitMQ、SQL Server Service Broker 这类纯队列 Provider 文档均明确要求 along with a distributed lock manager务必再注册一个UseDistributedLockManager/UseRedisLocking/UseAwsDynamoLocking等否则无法避免并发冲突。前缀prefix用于环境隔离Redis 队列键为{prefix}-workflows等、Redis 锁资源为{prefix}:{key}不同环境dev/staging/prod使用不同前缀即可复用同一基础设施同一集群内各节点必须使用相同前缀。同实例不重复执行WorkflowConsumer先取锁再执行、执行完释放锁配合共享队列同一时刻单实例单节点处理的语义在任意外部锁实现下保持一致DistributedLockProviderTests 验证了三个核心契约。事件订阅与生命周期事件集群模式下如需跨节点广播生命周期事件如 workflow 完成/出错可额外配置事件总线Redis EventHub、Azure Service Bus、AWS Kinesis保证各节点的OnLifeCycleEvent都能收到通知。定时/延迟执行的推进WorkflowConsumer对NextExecution的处理依赖Options.PollInterval默认 10s与持久化层的SupportsScheduledCommands能力——支持命令调度的存储如 SQL 系将写入ScheduledCommand否则由RunnablePoller轮询推进选型持久化时应关注此项差异。七、延伸阅读官方多节点文档本文所依据的 docs/multi-node-clusters.md接口契约IQueueProvider.cs、IDistributedLockProvider.cs默认实现SingleNodeQueueProvider.cs、SingleNodeLockProvider.cs消费与执行QueueConsumer.cs、WorkflowConsumer.cs、WorkflowHost.cs各 Provider 的 README 与扩展方法Redis、RabbitMQ、Azure、AWS、SQL Server锁契约测试DistributedLockProviderTests.cs配置选项WorkflowOptions.cs以上即为一篇完整的技术文章。我已基于 docs/multi-node-clusters.md 作为骨架结合仓库源码WorkflowOptions 默认值、SingleNodeQueueProvider/SingleNodeLockProvider、IQueueProvider/IDistributedLockProvider 接口、WorkflowConsumer 锁处理链路、QueueConsumer 并发控制、各 Provider 的 ServiceCollectionExtensions 与 README、DistributedLockProviderTests 测试基类进行了深度扩充覆盖了文档中列出的全部 ProviderAzure Storage Queues、Redis、RabbitMQ、AWS SQS 队列Azure Blob Leases、Redis、AWS DynamoDB 锁并补充了 SQL Server Service Broker 队列 Provider仓库内实际存在。全文为简体中文标题自拟所有链接均转换为仓库根目录相对路径未插入图片图片列表中无与该主题强相关的架构图/截屏无任何外部链接与营销性表述。赞分享后端工作流自动化流程编排【免费下载链接】workflow-coreLightweight workflow engine for .NET Standard项目地址https://gitcode.com/gh_mirrors/wo/workflow-core点击查看免费下载相关推荐HEIF vs JPEG为什么HEIF是更好的图像存储选择终极指南解析HEIF vs JPEG为什么HEIF是更好的图像存储选择终极指南解析 在当今数字图像存储领域 HEIFHigh Efficiency Image FiSGLang多节点部署分布式集群的配置与管理SGLang多节点部署分布式集群的配置与管理 概述 在大模型推理场景中当单个节点的计算资源无法满足超大规模模型如Llama 3.1 405B、DeepSe模型推理服务推理引擎人工智能大模型本地部署多模态CrowdSec 分布式系统部署指南多节点集群部署步骤CrowdSec 分布式系统部署指南多节点集群部署步骤 随着网络攻击日益复杂化单一节点的安全防护系统已难以应对大规模、分布式的威胁。CrowdSec 作为一网络安全应用安全WAFIDSIPS上一篇技术突破CogVLM如何重新定义视觉语言模型的能力边界下一篇解决SLIM容器文件系统挂载权限从错误排查到最佳实践创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
阅读完成 · 觉得有帮助?