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

基于IEnlistmentNotification实现Azure Blob通用两阶段提交类

通用Azure Blob事务操作实现方案

核心思路

通过抽象基类+子类的方式封装不同Blob操作的业务逻辑与补偿(回滚)逻辑,让资源管理器只依赖抽象契约,从而支持任意Blob操作的事务集成。具体设计:

  • 定义抽象基类统一操作接口,包含执行操作和回滚操作的契约
  • 每个Blob操作(上传、删除、设置元数据)对应一个子类,实现自身的业务逻辑与回滚补偿逻辑
  • 资源管理器维护抽象操作的列表,无需关心具体操作类型

代码实现

1. 抽象事务操作基类

public abstract class BlobTransactionalOperation : IDisposable
{
    protected readonly BlobClient BlobClient;
    private bool _disposedValue;

    protected BlobTransactionalOperation(BlobContainerClient containerClient, string blobName)
    {
        BlobClient = containerClient.GetBlobClient(blobName);
    }

    // 执行核心业务操作
    public abstract Task DoWork();

    // 执行回滚补偿操作
    public abstract Task RollBack();

    public void Dispose() => Dispose(true);

    protected virtual void Dispose(bool disposing)
    {
        if (_disposedValue) return;
        if (disposing)
        {
            // 子类可在此释放特定资源
        }
        _disposedValue = true;
    }

    ~BlobTransactionalOperation() => Dispose(false);
}

2. 具体操作子类

上传Blob操作

public class BlobUploadOperation : BlobTransactionalOperation
{
    private readonly Stream _content;

    public BlobUploadOperation(BlobContainerClient containerClient, string blobName, Stream content)
        : base(containerClient, blobName)
    {
        _content = content;
        _content.Position = 0;
    }

    public override async Task DoWork()
    {
        await BlobClient.UploadAsync(_content, overwrite: true);
    }

    public override async Task RollBack()
    {
        // 补偿逻辑:删除已上传的Blob
        await BlobClient.DeleteIfExistsAsync(Azure.Storage.Blobs.Models.DeleteSnapshotsOption.IncludeSnapshots);
    }
}

删除Blob操作

public class BlobDeleteOperation : BlobTransactionalOperation
{
    private Stream _backupContent;
    private BlobProperties _originalProperties;

    public BlobDeleteOperation(BlobContainerClient containerClient, string blobName)
        : base(containerClient, blobName)
    {
    }

    public override async Task DoWork()
    {
        // 先备份原Blob内容和属性,用于回滚
        if (await BlobClient.ExistsAsync())
        {
            _backupContent = new MemoryStream();
            await BlobClient.DownloadToAsync(_backupContent);
            _backupContent.Position = 0;
            _originalProperties = (await BlobClient.GetPropertiesAsync()).Value;
        }

        // 执行删除操作
        await BlobClient.DeleteAsync(Azure.Storage.Blobs.Models.DeleteSnapshotsOption.IncludeSnapshots);
    }

    public override async Task RollBack()
    {
        // 补偿逻辑:恢复备份的Blob
        if (_backupContent != null)
        {
            _backupContent.Position = 0;
            var uploadOptions = new BlobUploadOptions
            {
                Properties = _originalProperties
            };
            await BlobClient.UploadAsync(_backupContent, uploadOptions);
        }
    }

    protected override void Dispose(bool disposing)
    {
        base.Dispose(disposing);
        if (disposing)
        {
            _backupContent?.Dispose();
        }
    }
}

设置Blob元数据操作

public class BlobSetMetadataOperation : BlobTransactionalOperation
{
    private readonly IDictionary<string, string> _newMetadata;
    private IDictionary<string, string> _originalMetadata;

    public BlobSetMetadataOperation(BlobContainerClient containerClient, string blobName, IDictionary<string, string> newMetadata)
        : base(containerClient, blobName)
    {
        _newMetadata = newMetadata;
    }

    public override async Task DoWork()
    {
        // 先备份原元数据
        if (await BlobClient.ExistsAsync())
        {
            var properties = await BlobClient.GetPropertiesAsync();
            _originalMetadata = new Dictionary<string, string>(properties.Value.Metadata);
        }

        // 设置新元数据
        await BlobClient.SetMetadataAsync(_newMetadata);
    }

    public override async Task RollBack()
    {
        // 补偿逻辑:恢复原元数据
        if (_originalMetadata != null)
        {
            await BlobClient.SetMetadataAsync(_originalMetadata);
        }
    }
}

3. 调整后的资源管理器

public class AzureBlobStorageResourceManager : IEnlistmentNotification, IDisposable
{
    private List<BlobTransactionalOperation> _operations;
    private bool _disposedValue;

    public void EnlistOperation(BlobTransactionalOperation operation)
    {
        if (_operations == null)
        {
            var currentTransaction = Transaction.Current;
            currentTransaction?.EnlistVolatile(this, EnlistmentOptions.None);
            _operations = new List<BlobTransactionalOperation>();
        }
        _operations.Add(operation);
    }

    public void Commit(Enlistment enlistment)
    {
        foreach (var operation in _operations)
        {
            operation.Dispose();
        }
        enlistment.Done();
    }

    public void InDoubt(Enlistment enlistment)
    {
        foreach (var operation in _operations)
        {
            operation.RollBack().ConfigureAwait(false).GetAwaiter().GetResult();
        }
        enlistment.Done();
    }

    public void Prepare(PreparingEnlistment preparingEnlistment)
    {
        try
        {
            foreach (var operation in _operations)
            {
                operation.DoWork().ConfigureAwait(false).GetAwaiter().GetResult();
            }
            preparingEnlistment.Prepared();
        }
        catch
        {
            preparingEnlistment.ForceRollback();
        }
    }

    public void Rollback(Enlistment enlistment)
    {
        foreach (var operation in _operations)
        {
            operation.RollBack().ConfigureAwait(false).GetAwaiter().GetResult();
        }
        enlistment.Done();
    }

    public void Dispose() => Dispose(true);

    protected virtual void Dispose(bool disposing)
    {
        if (_disposedValue) return;

        if (disposing)
        {
            foreach (var operation in _operations)
                operation.Dispose();
        }

        _disposedValue = true;
    }

    ~AzureBlobStorageResourceManager() => Dispose(false);
}

使用示例

using (var scope = new TransactionScope(TransactionScopeAsyncFlowOption.Enabled))
{
    var containerClient = new BlobContainerClient("<Blob连接字符串>", "<容器名称>");
    var blobResourceManager = new AzureBlobStorageResourceManager();

    // 1. 注册Blob上传操作
    using var uploadStream = File.OpenRead("test.txt");
    var uploadOp = new BlobUploadOperation(containerClient, "test.txt", uploadStream);
    blobResourceManager.EnlistOperation(uploadOp);

    // 2. 注册Blob元数据设置操作
    var metadataOp = new BlobSetMetadataOperation(containerClient, "test.txt", new Dictionary<string, string> { { "Category", "Document" } });
    blobResourceManager.EnlistOperation(metadataOp);

    // 3. 执行Azure SQL事务操作(示例)
    using var sqlConn = new SqlConnection("<SQL连接字符串>");
    await sqlConn.OpenAsync();
    using var cmd = sqlConn.CreateCommand();
    cmd.CommandText = "INSERT INTO Documents (Name) VALUES ('test.txt')";
    await cmd.ExecuteNonQueryAsync();

    // 提交事务
    scope.Complete();
}

内容的提问来源于stack exchange,提问作者Kalaivanan2k4

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 09:25:40