基于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
相关产品推荐
相关产品推荐

