Elsa.Client的IWorkflowDefinitionsApi是否支持数据库事务及并发问题咨询
Elsa 2.14.1 工作流定义并发更新丢失IsLatest=true记录问题及解决方案
问题详情
- 使用版本:Elsa 2.14.1
- 操作场景:通过
Elsa.Client.Services.IWorkflowDefinitionsApi的SaveAsync端点(/v1/workflow-definitions/{workflowDefinitionVersionId})存储动态生成的工作流,数据库采用Docker容器部署的PostgreSQL - 核心问题:当2个并发用户更新同一ID的工作流定义时,会出现特定工作流无
IsLatest列设为true的记录的情况 - 推测原因:Elsa Server保存工作流定义时,未用事务包裹“将旧版本IsLatest设为false”和“插入/更新新版本并设IsLatest=true”这两步操作,并发下导致中间状态丢失
复现代码
using Elsa.Client; using Elsa.Client.Extensions; using Elsa.Client.Models; using Microsoft.Extensions.DependencyInjection; ServiceCollection serviceCollection = new(); serviceCollection.AddElsaClient( // 连接Docker容器部署的Elsa服务 options => options.ServerUrl = new Uri("http://localhost:13000")); var sp = serviceCollection.BuildServiceProvider(); var tasks = Enumerable .Range(1, 2) .Select(i => { var threadNo = i.ToString(); return new TaskFactory( TaskCreationOptions.LongRunning, TaskContinuationOptions.None) .StartNew(async () => await RunSaveLoop(threadNo)); }) .Select(t=> t.Result) .ToArray(); await Task.WhenAll(tasks); return 0; async Task RunSaveLoop(string id) { await Task.Yield(); while (true) { await SaveWorkflow(id); } } async Task SaveWorkflow( string id) { Console.WriteLine($"Saving workflow in thread {id}"); var request = new SaveWorkflowDefinitionRequest { WorkflowDefinitionId = "test definition", Name = "test definition", DisplayName = "test definition", Publish = true, Activities = new List<ActivityDefinition>(), Connections = new List<ConnectionDefinition>(), }; var elsaClient = sp.GetRequiredService<IElsaClient>(); try { var response = await elsaClient .WorkflowDefinitions .SaveAsync(request); Console.WriteLine($"Saved workflow in thread {id}"); } catch (Exception e) { Console.WriteLine($"Thread {id} failed: {e.Message}"); throw; } }
解决方案
1. 自定义带事务的工作流定义存储服务
Elsa支持替换默认的IWorkflowDefinitionStore实现,我们可以创建一个事务包装类,将保存操作包裹在事务中:
public class TransactionalWorkflowDefinitionStore : IWorkflowDefinitionStore { private readonly IWorkflowDefinitionStore _innerStore; private readonly IDbContextTransactionManager _transactionManager; public TransactionalWorkflowDefinitionStore(IWorkflowDefinitionStore innerStore, IDbContextTransactionManager transactionManager) { _innerStore = innerStore; _transactionManager = transactionManager; } public async Task<WorkflowDefinition> SaveAsync(WorkflowDefinition definition, CancellationToken cancellationToken = default) { using var transaction = await _transactionManager.BeginTransactionAsync(cancellationToken); try { var result = await _innerStore.SaveAsync(definition, cancellationToken); await transaction.CommitAsync(cancellationToken); return result; } catch { await transaction.RollbackAsync(cancellationToken); throw; } } // 其他接口方法直接委托给内部存储实现 public Task<WorkflowDefinition?> FindByIdAsync(string id, VersionOptions versionOptions, CancellationToken cancellationToken = default) => _innerStore.FindByIdAsync(id, versionOptions, cancellationToken); public Task<WorkflowDefinition?> FindByDefinitionIdAsync(string definitionId, VersionOptions versionOptions, CancellationToken cancellationToken = default) => _innerStore.FindByDefinitionIdAsync(definitionId, versionOptions, cancellationToken); // 剩余接口方法同理实现... }
然后在Elsa启动配置中替换服务:
services.AddElsa(elsa => { elsa.ReplaceService<IWorkflowDefinitionStore, TransactionalWorkflowDefinitionStore>(ServiceLifetime.Scoped); });
2. 利用工作流保存钩子注入事务逻辑
通过Elsa的IWorkflowDefinitionSavedHandler钩子,在保存前后管理事务:
public class TransactionalWorkflowDefinitionSavedHandler : IWorkflowDefinitionSavedHandler { private readonly IDbContextTransactionManager _transactionManager; private IDbContextTransaction? _transaction; public TransactionalWorkflowDefinitionSavedHandler(IDbContextTransactionManager transactionManager) { _transactionManager = transactionManager; } public async Task HandleAsync(WorkflowDefinitionSavedContext context, CancellationToken cancellationToken) { if (context.IsPreSave) { // 保存前开启事务 _transaction = await _transactionManager.BeginTransactionAsync(cancellationToken); } else { // 保存完成后提交事务 if (_transaction != null) { await _transaction.CommitAsync(cancellationToken); await _transaction.DisposeAsync(); } } } }
注册该钩子:
services.AddScoped<IWorkflowDefinitionSavedHandler, TransactionalWorkflowDefinitionSavedHandler>();
3. 数据库层面行级锁(辅助方案)
在PostgreSQL中,针对workflow_definitions表的workflow_definition_id字段添加行级锁,避免并发更新的竞态。可以修改Elsa底层数据访问的SQL语句,在更新旧版本IsLatest状态时使用SELECT ... FOR UPDATE锁定目标行,确保同一时间只有一个操作能修改该工作流的版本状态。
注意事项
- 事务粒度需控制在合理范围,避免过度锁表影响系统性能
- 替换存储服务时,需确保所有工作流定义的保存路径都使用了新的事务化服务
- 测试时需模拟高并发场景,验证事务是否彻底解决
IsLatest记录丢失问题
内容的提问来源于stack exchange,提问作者Andrey A
相关产品推荐
相关产品推荐

