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

如何基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 12:35:31