如何实现含DB写入、Blob上传、异步消息的可补偿事务与回滚
问题解答
关于TransactionScope方案的可行性
你最初用TransactionScope包装全流程的思路无法实现预期的自动回滚:
- EF Core搭配Npgsql确实支持环境事务,仅数据库写入操作会在异常时自动回滚,这部分你的理解没有问题。
- Azure Blob存储、Service Bus官方客户端原生没有实现System.Transactions的事务接入接口,不会感知TransactionScope的提交/回滚状态,只要对应API调用成功,操作就会立刻落地,不会等待
scope.Commit()执行,后续步骤失败时这两个操作不会自动回滚。 - 自行实现IEnlistmentNotification接口包装两个客户端做自定义资源管理器理论上可行,但生产环境容错性极差:进程崩溃、网络瞬断、事务状态不确定等问题都会导致脏数据残留,云场景下本身也不推荐用两阶段提交的分布式事务方案,维护成本极高。
各操作的补偿逻辑设计
- Blob存储的补偿逻辑你判断正确,回滚时直接调用
DeleteIfExistsAsync方法删除已上传的文件即可,注意做幂等判断,文件不存在时直接视为补偿成功,不要抛出异常阻断流程。 - Service Bus没有提供已投递到队列的消息撤回接口,不要尝试发完消息再删除,正确的设计是利用原生的延迟投递特性实现补偿:
- 发送消息时设置
ScheduledEnqueueTimeUtc属性,给消息预留足够覆盖全流程执行时长的延迟窗口(比如当前时间加30秒),延迟到期前消费者完全拉取不到这条消息 - 流程失败需要回滚时,调用
CancelScheduledMessageAsync传入发送时返回的消息序列号,就能直接取消这条延迟消息的投递 - 全流程执行成功后,可以直接取消延迟消息、立刻发送一条无延迟的消息供消费者消费,也可以等延迟窗口到期自动投递。
- 发送消息时设置
- 数据库操作直接用EF Core自带的本地事务即可,不需要引入分布式事务能力。
推荐实现流程
你判断Saga模式不适用是完全正确的:Saga是面向跨微服务、跨进程长事务的一致性方案,你当前所有操作都在单服务进程内,流程短、补偿逻辑简单,不需要引入Saga的编排/协同复杂度,直接按「先执行无状态/易补偿操作、最后提交数据库事务」的顺序做失败补偿即可,流程如下:
- 第一步执行Blob上传,这步失败直接抛出错误即可,没有任何持久化数据需要回滚
- 第二步发送带延迟属性的Service Bus消息,记录返回的消息序列号,这步如果失败,直接删除第一步上传的Blob后抛出错误
- 第三步开启EF Core本地事务,写入实体数据并保存,这步如果失败,依次执行删除Blob、取消延迟消息的补偿逻辑,回滚数据库事务后抛出错误
- 所有步骤执行无异常后,提交数据库事务,再触发即时消息通知,流程结束
参考实现代码如下:
long? scheduledMessageSeq = null; using var dbTransaction = await _dbContext.Database.BeginTransactionAsync(); try { // 上传Blob var blobClient = _blobContainerClient.GetBlobClient(entity.FileIdentifier); await blobClient.UploadAsync(entity.FileStream); // 发送30秒延迟的调度消息 var queueMessage = new ServiceBusMessage(BinaryData.FromObjectAsJson(new { entity.Id })) { ScheduledEnqueueTimeUtc = DateTimeOffset.UtcNow.AddSeconds(30) }; scheduledMessageSeq = await _serviceBusSender.ScheduleMessageAsync( queueMessage, queueMessage.ScheduledEnqueueTimeUtc ); // 写入数据库 await _dbContext.Set<FileEntity>().AddAsync(entity); await _dbContext.SaveChangesAsync(); // 提交数据库事务 await dbTransaction.CommitAsync(); // 全流程成功,取消延迟消息,发送即时通知 if (scheduledMessageSeq.HasValue) { await _serviceBusSender.CancelScheduledMessageAsync(scheduledMessageSeq.Value); await _serviceBusSender.SendMessageAsync( new ServiceBusMessage(BinaryData.FromObjectAsJson(new { entity.Id })) ); } } catch (Exception) { await dbTransaction.RollbackAsync(); // 异步执行补偿,不阻塞异常抛出 _ = Task.Run(async () => { try { await _blobContainerClient.GetBlobClient(entity.FileIdentifier).DeleteIfExistsAsync(); } catch { /* 记录日志即可,后续可通过定时巡检清理残留Blob */ } if (scheduledMessageSeq.HasValue) { try { await _serviceBusSender.CancelScheduledMessageAsync(scheduledMessageSeq.Value); } catch { /* 记录日志即可,消费者侧做幂等校验,查询不到对应实体直接丢弃消息 */ } } }); throw; }
兜底设计:就算补偿操作因为网络等问题执行失败,也不会产生严重一致性问题:残留Blob可以通过定时任务对比数据库记录清理,漏取消的延迟消息投递后,消费者只要先查数据库中对应实体是否存在,不存在直接丢弃消息即可,容错成本远低于分布式两阶段提交。
内容的提问来源于stack exchange,提问作者Matthias Güntert
相关产品推荐
相关产品推荐

