基于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
相关产品推荐
相关产品推荐

