如何基于MassTransit事务性发件箱设计批量邮件发送系统?
事务性发件箱实现:避免重复邮件的正确设计
原始代码问题
现有消费逻辑直接在数据库保存后发送邮件,存在重复发送风险:如果邮件发送失败后触发重试消费,会重复执行发送操作。
private readonly IEmailSender _emailSender; public async Task Consume(ConsumeContext<UserCreated> context) // UserCreated 来自外部系统事件 { await _dbContext.User.AddAsync(new UserAggregate { UserName = "kiddo" }); EmailTemplate[] emailTemplates = CreateEmailTemplatePerUser(); await _dbContext.SaveChangesAsync(); await _emailSender.Send(emailTemplates); }
初步重构的问题
为引入发件箱模式,替换IEmailSender为ISendEndpointProvider,但代码顺序错误,无法保证事务一致性:
private readonly ISendEndpointProvider _sendEndpointProvider; public async Task Consume(ConsumeContext<UserCreated> context) { await _dbContext.User.AddAsync(new UserAggregate { UserName = "kiddo" }); EmailCommands[] emailCommands = CreateEmailCommandPerUser(); await _sendEndpointProvider.Send(emailCommands); await _dbContext.SaveChangesAsync(); }
这段代码先发送消息再保存数据库,若数据库保存失败,已发送的消息会导致重复执行。
需求明确
我们需要:
- 用发件箱生成独立的单封邮件发送命令(实际场景会按100封批量处理,示例已简化)
- 使用Mongo事务性发件箱,将消息持久化到
outbox.messages集合,而非通过ConsumeContext立即发送 - 避免创建新作用域脱离当前
ConsumeContext,防止作用域过滤器产生副作用
正确实现方案
1. 配置Mongo事务性发件箱
先在MassTransit配置中启用Mongo发件箱,确保消息会被持久化到指定集合:
services.AddMassTransit(x => { x.SetKebabCaseEndpointNameFormatter(); x.UsingRabbitMq((context, cfg) => { cfg.Host("localhost"); // 配置事务性发件箱 cfg.UseMongoDbOutbox(context, "YourMongoConnectionString") .UseCollectionName("outbox.messages") .UseBusOutbox(); }); });
2. 调整消费逻辑(事务内原子操作)
修改消费方法,将数据库操作和发件箱消息持久化放在同一个Mongo事务中,确保两者原子性:
private readonly ISendEndpointProvider _sendEndpointProvider; private readonly IMongoDatabase _database; public async Task Consume(ConsumeContext<UserCreated> context) { using var session = await _database.Client.StartSessionAsync(); session.StartTransaction(); try { // 事务内添加用户数据 var userCollection = _database.GetCollection<UserAggregate>("users"); await userCollection.InsertOneAsync(session, new UserAggregate { UserName = "kiddo" }); // 生成邮件命令 EmailCommands[] emailCommands = CreateEmailCommandPerUser(); // 通过ConsumeContext发送,自动路由到发件箱 foreach (var command in emailCommands) { await context.Send(command); } // 提交事务:同时保存数据库变更和发件箱消息 await session.CommitTransactionAsync(); } catch (Exception ex) { // 回滚事务:所有操作撤销 await session.AbortTransactionAsync(); throw; } }
3. 核心注意事项
- 事务原子性:数据库操作和发件箱消息必须在同一个事务中,要么全部成功,要么全部回滚
- 使用ConsumeContext发送:通过
context.Send()替代直接调用_sendEndpointProvider.Send(),MassTransit会自动将消息存入发件箱,而非立即投递 - 发件箱后台处理:事务提交后,MassTransit的发件箱后台进程会自动处理
outbox.messages中的消息,确保每个消息只被处理一次 - 批量处理实现:批量发送逻辑应放在
EmailCommands的消费者中,根据第三方API限制进行批量打包,当前消费方法只负责生成独立命令,保持职责单一
内容的提问来源于stack exchange,提问作者Ilias Dewachter
相关产品推荐
相关产品推荐

