如何在MassTransit中结合RavenDb实现持久化存储与Transactional Outbox?
实现RavenDB + MassTransit Transactional Outbox的建议与方向
核心思路
Transactional Outbox的核心是在同一事务中原子性保存业务数据和待发送消息,之后通过异步可靠的方式投递消息,并自动清理已处理的记录。结合RavenDB的ACID事务特性、内置Change Feed推送能力,以及MassTransit的消息投递能力,可以实现轻量且可靠的Outbox机制。
具体实现步骤
1. 定义Outbox消息实体
在RavenDB中创建专门的Outbox消息文档,用于存储待发送的消息元数据和内容:
public class OutboxMessage { public string Id { get; set; } public byte[] MessageBody { get; set; } public string MessageType { get; set; } // 消息类型的程序集限定名,用于反序列化 public OutboxMessageStatus Status { get; set; } public DateTime CreatedAt { get; set; } public int RetryCount { get; set; } public DateTime? ProcessedAt { get; set; } public DateTime? ExpiresAt { get; set; } // 用于RavenDB自动过期清理 } public enum OutboxMessageStatus { Pending, // 待处理 Processing, // 处理中(防止重复投递) Completed, // 投递成功 Failed // 重试耗尽后失败 }
2. 业务事务中写入Outbox
在业务操作的RavenDB会话事务中,同时保存业务数据和Outbox消息,确保两者原子性提交:
using var session = _documentStore.OpenAsyncSession(); using var transaction = await session.Advanced.StartTransactionAsync(); try { // 1. 保存业务数据(示例:订单创建) var order = new Order { Id = "orders/123", Status = OrderStatus.Created }; await session.StoreAsync(order); // 2. 序列化消息并写入Outbox var orderCreatedEvent = new OrderCreated { OrderId = order.Id }; var messageBody = JsonSerializer.SerializeToUtf8Bytes(orderCreatedEvent); var outboxMsg = new OutboxMessage { MessageBody = messageBody, MessageType = typeof(OrderCreated).AssemblyQualifiedName, Status = OutboxMessageStatus.Pending, CreatedAt = DateTime.UtcNow, RetryCount = 0 }; await session.StoreAsync(outboxMsg); // 3. 提交事务 await session.SaveChangesAsync(); await transaction.CommitAsync(); } catch { await transaction.RollbackAsync(); throw; }
3. 基于RavenDB Change Feed实现消息投递
利用RavenDB的Change Feed监听新的Outbox消息,实时触发MassTransit的消息投递:
// 注册Change Feed观察者(建议作为HostedService运行) public class OutboxProcessorHostedService : IHostedService { private readonly IBus _bus; private readonly IDocumentStore _documentStore; private IDisposable _changeSubscription; public OutboxProcessorHostedService(IBus bus, IDocumentStore documentStore) { _bus = bus; _documentStore = documentStore; } public async Task StartAsync(CancellationToken cancellationToken) { var observer = new OutboxMessageObserver(_bus, _documentStore); _changeSubscription = await _documentStore.Changes() .ForDocumentsStartingWith("OutboxMessage/") // 匹配Outbox消息文档ID前缀 .SubscribeAsync(observer, cancellationToken); } public Task StopAsync(CancellationToken cancellationToken) { _changeSubscription?.Dispose(); return Task.CompletedTask; } } // Change Feed观察者实现 public class OutboxMessageObserver : IObserver<DocumentChange> { private readonly IBus _bus; private readonly IDocumentStore _documentStore; private readonly int _maxRetryCount = 5; public OutboxMessageObserver(IBus bus, IDocumentStore documentStore) { _bus = bus; _documentStore = documentStore; } public async void OnNext(DocumentChange change) { if (change.Type != DocumentChangeType.Put) return; using var session = _documentStore.OpenAsyncSession(); var outboxMsg = await session.LoadAsync<OutboxMessage>(change.Id); if (outboxMsg?.Status != OutboxMessageStatus.Pending) return; // 先标记为处理中,避免重复投递 outboxMsg.Status = OutboxMessageStatus.Processing; await session.SaveChangesAsync(); try { // 反序列化消息并通过MassTransit发布 var messageType = Type.GetType(outboxMsg.MessageType); var message = JsonSerializer.Deserialize(outboxMsg.MessageBody, messageType); await _bus.Publish(message, messageType); // 标记为完成,并设置过期时间(7天后自动清理) outboxMsg.Status = OutboxMessageStatus.Completed; outboxMsg.ProcessedAt = DateTime.UtcNow; outboxMsg.ExpiresAt = DateTime.UtcNow.AddDays(7); } catch (Exception ex) { outboxMsg.RetryCount++; if (outboxMsg.RetryCount >= _maxRetryCount) { outboxMsg.Status = OutboxMessageStatus.Failed; outboxMsg.ExpiresAt = DateTime.UtcNow.AddDays(30); // 失败消息保留更久用于排查 // 这里可以添加告警逻辑,比如写入错误日志或通知运维 } else { outboxMsg.Status = OutboxMessageStatus.Pending; // 允许后续重试 } } finally { await session.SaveChangesAsync(); } } public void OnError(Exception error) { // 处理Change Feed监听错误,比如记录日志 } public void OnCompleted() { // 监听终止处理 } }
4. 配置RavenDB自动清理过期Outbox消息
启用RavenDB的文档过期功能,自动清理已完成或失败的Outbox消息:
// 在应用启动时配置过期策略 await _documentStore.Maintenance.SendAsync(new ConfigureExpirationOperation { Database = _documentStore.Database, Settings = new ExpirationSettings { Enabled = true, DeleteFrequency = TimeSpan.FromHours(1), // 每小时清理一次过期文档 RemoveFrequency = TimeSpan.FromHours(6) } });
关键注意事项
- 事务一致性:必须确保业务数据和Outbox消息在同一个RavenDB会话事务中提交,避免数据与消息不一致。
- 幂等性:消息消费者需要实现幂等处理,防止因重试导致的重复业务操作。
- 重试机制:设置合理的重试次数上限,避免无效重试占用资源;失败消息保留足够时间用于问题排查。
- 并发控制:通过
Processing状态标记,避免同一消息被多个处理器同时处理。
内容的提问来源于stack exchange,提问作者James Hancock
相关产品推荐
相关产品推荐

