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

基于MassTransit+Azure Service Bus,数据库更新后调用外部API发邮件咨询

实现步骤及代码示例

1. 定义事件消息类(若未定义)

首先确保UpdatedEventMessage已正确定义,用于传递注册编号:

public class UpdatedEventMessage
{
    public string RegistNumber { get; set; }
}

2. 创建事件消费者

实现MassTransit的IConsumer<UpdatedEventMessage>,整合数据库查询、外部API调用、邮件发送逻辑:

public class UpdatedEventConsumer : IConsumer<UpdatedEventMessage>
{
    private readonly AppDbContext _dbContext;
    private readonly ExternalUserApiClient _userApiClient;
    private readonly EmailService _emailService;

    // 通过依赖注入获取所需服务
    public UpdatedEventConsumer(AppDbContext dbContext, ExternalUserApiClient userApiClient, EmailService emailService)
    {
        _dbContext = dbContext;
        _userApiClient = userApiClient;
        _emailService = emailService;
    }

    public async Task Consume(ConsumeContext<UpdatedEventMessage> context)
    {
        var registNumber = context.Message.RegistNumber;

        // 步骤1:从本地数据库查询用户ID
        var userId = await _dbContext.Registrations
            .Where(r => r.RegistNumber == registNumber)
            .Select(r => r.UserId)
            .FirstOrDefaultAsync();

        if (userId == Guid.Empty)
        {
            await context.LogError($"未找到RegistNumber为{registNumber}的用户ID");
            return;
        }

        // 步骤2:调用外部API获取用户邮箱和称呼
        UserInfo userInfo;
        try
        {
            userInfo = await _userApiClient.GetUserInfoById(userId);
        }
        catch (HttpRequestException ex)
        {
            await context.LogError(ex, "调用外部用户API失败,用户ID:{UserId}", userId);
            return;
        }

        if (string.IsNullOrWhiteSpace(userInfo.Email))
        {
            await context.LogError("用户ID:{UserId}未配置邮箱地址", userId);
            return;
        }

        // 步骤3:发送通知邮件
        try
        {
            await _emailService.SendUpdateNotificationEmail(userInfo);
            await context.LogInformation("已成功向{Email}发送更新通知邮件", userInfo.Email);
        }
        catch (Exception ex)
        {
            await context.LogError(ex, "向{Email}发送邮件失败", userInfo.Email);
        }
    }
}

3. 实现本地数据库查询

使用EF Core实现从本地数据库获取用户ID的逻辑(假设你已有Registrations实体):

// 数据库上下文示例
public class AppDbContext : DbContext
{
    public AppDbContext(DbContextOptions<AppDbContext> options) : base(options) { }

    public DbSet<Registration> Registrations { get; set; }
}

public class Registration
{
    public Guid Id { get; set; }
    public string RegistNumber { get; set; }
    public Guid UserId { get; set; }
    // 其他字段...
}

4. 实现外部API调用客户端

用HttpClient封装外部API调用,获取用户信息:

public class ExternalUserApiClient
{
    private readonly HttpClient _httpClient;

    public ExternalUserApiClient(HttpClient httpClient)
    {
        _httpClient = httpClient;
        // 从配置文件读取API地址,避免硬编码
        _httpClient.BaseAddress = new Uri(Environment.GetEnvironmentVariable("ExternalUserApiBaseUrl"));
    }

    public async Task<UserInfo> GetUserInfoById(Guid userId)
    {
        var response = await _httpClient.GetAsync($"api/users/{userId}");
        response.EnsureSuccessStatusCode(); // 非2xx状态码抛出异常

        return await response.Content.ReadFromJsonAsync<UserInfo>();
    }
}

// 外部API返回的用户信息模型
public class UserInfo
{
    public string Email { get; set; }
    public string FirstName { get; set; }
    public string LastName { get; set; }
}

5. 实现邮件发送服务

用MailKit实现邮件发送(也可替换为SendGrid等服务):

public class EmailService
{
    private readonly IConfiguration _configuration;

    public EmailService(IConfiguration configuration)
    {
        _configuration = configuration;
    }

    public async Task SendUpdateNotificationEmail(UserInfo userInfo)
    {
        using var smtpClient = new SmtpClient();
        // 从配置文件读取SMTP参数
        await smtpClient.ConnectAsync(
            _configuration["Smtp:Host"],
            int.Parse(_configuration["Smtp:Port"]),
            SecureSocketOptions.StartTls);
        
        await smtpClient.AuthenticateAsync(
            _configuration["Smtp:Username"],
            _configuration["Smtp:Password"]);

        var message = new MimeMessage();
        message.From.Add(new MailboxAddress("系统通知", _configuration["Smtp:FromEmail"]));
        message.To.Add(new MailboxAddress($"{userInfo.FirstName} {userInfo.LastName}", userInfo.Email));
        message.Subject = "您的注册信息已更新";

        message.Body = new TextPart("html")
        {
            Text = $@"
                <h3>您好,{userInfo.FirstName} {userInfo.LastName}:</h3>
                <p>您的注册信息已成功完成更新,如有疑问请联系客服。</p>
            "
        };

        await smtpClient.SendAsync(message);
        await smtpClient.DisconnectAsync(true);
    }
}

6. 配置MassTransit与Azure Service Bus

在Program.cs中注册所有服务并配置MassTransit:

var builder = WebApplication.CreateBuilder(args);

// 注册数据库上下文
builder.Services.AddDbContext<AppDbContext>(options =>
    options.UseSqlServer(builder.Configuration.GetConnectionString("DefaultConnection")));

// 注册外部API客户端
builder.Services.AddHttpClient<ExternalUserApiClient>();

// 注册邮件服务
builder.Services.AddTransient<EmailService>();

// 配置MassTransit和Azure Service Bus
builder.Services.AddMassTransit(x =>
{
    // 注册事件消费者
    x.AddConsumer<UpdatedEventConsumer>();

    x.UsingAzureServiceBus((context, cfg) =>
    {
        // 从配置文件读取Service Bus连接字符串
        cfg.Host(builder.Configuration["AzureServiceBus:ConnectionString"]);

        // 配置消费者接收端点
        cfg.ReceiveEndpoint("updated-event-email-handler", e =>
        {
            e.ConfigureConsumer<UpdatedEventConsumer>(context);

            // 配置重试策略:API调用或邮件发送失败时重试3次,间隔5秒
            e.UseMessageRetry(r => r.Interval(3, TimeSpan.FromSeconds(5)));

            // 配置死信队列:重试失败后将消息移入死信队列,避免重复消费
            e.DeadLetterQueueName = "updated-event-dead-letter";
        });
    });
});

var app = builder.Build();

// 其他中间件配置...

app.Run();

关键注意事项

  • 错误处理:通过日志记录失败场景,利用MassTransit的重试和死信队列机制避免消息丢失
  • 配置管理:所有敏感信息(连接字符串、API地址、SMTP参数)均从配置文件或环境变量读取,禁止硬编码
  • 依赖注入:所有服务通过DI容器注入,保证代码可测试性和扩展性

内容的提问来源于stack exchange,提问作者Dhanyanjay Singh

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 12:16:07