You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.12 17:07:14