基于最终一致性的一对一关系:CreateCustomer用例实现问询
微服务最终一致性实践:CreateCustomer用例实现方案
问题1:CreateCustomer用例应由哪个微服务负责实现?
你的判断完全正确,microservice_B应该主导这个用例。因为CUSTOMER是microservice_B的核心聚合根,CreateCustomer是围绕CUSTOMER生命周期展开的业务操作,由归属的微服务负责更符合领域驱动设计(DDD)的聚合根职责边界原则。
问题2:CreateCustomer用例如何实现?有哪些最佳实践?
结合最终一致性原则和.NET技术栈生态,针对你提出的两个方案,以下是优化后的落地实践:
方案一:事件驱动的最终一致性(推荐)
异步事件驱动是实现跨微服务最终一致性的标准方案,需补充状态跟踪和重试/补偿机制解决中间状态卡顿问题,避免业务流程中断。
具体流程
- microservice_B接收CreateCustomer请求,创建**状态为"待关联账户"**的CUSTOMER记录(用0/Guid.Empty等标记代替null,避免违反"每个CUSTOMER必须关联ACCOUNT"的业务规则,仅标识关联未完成)
- microservice_B在本地事务内发布
CustomerCreatedPendingAccount事件(包含CustomerId、CustomerName) - microservice_A监听事件,通过自身聚合工厂创建符合业务规则的ACCOUNT,生成AccountId
- microservice_A发布
AccountCreated事件(包含AccountId、AccountName、关联的CustomerId) - microservice_B监听事件,更新对应CUSTOMER的AccountId,并将状态改为"正常"
- 配置死信队列和定时补偿任务:若某一步事件处理失败,将事件转入死信队列,通过定时任务重试;若重试仍失败,触发人工介入排查。
.NET技术栈实现示例
使用MassTransit(.NET成熟消息框架)实现事件总线,结合EF Core保证本地事务一致性,内置重试和死信机制。
1. microservice_B:发布创建待关联Customer事件
// 定义事件契约 public record CustomerCreatedPendingAccount(Guid CustomerId, string CustomerName); // Customer业务服务 public class CustomerService { private readonly AppDbContext _dbContext; private readonly IPublishEndpoint _publishEndpoint; public CustomerService(AppDbContext dbContext, IPublishEndpoint publishEndpoint) { _dbContext = dbContext; _publishEndpoint = publishEndpoint; } public async Task CreateCustomerAsync(string customerName) { // 用聚合工厂创建待关联状态的Customer(符合microservice_B业务规则) var customer = CustomerFactory.CreatePendingAccountCustomer(customerName); await _dbContext.Customers.AddAsync(customer); // 本地事务:保存Customer与发布事件原子执行 await _dbContext.Database.BeginTransactionAsync(); try { await _dbContext.SaveChangesAsync(); await _publishEndpoint.Publish(new CustomerCreatedPendingAccount(customer.CustomerId, customer.CustomerName)); await _dbContext.Database.CommitTransactionAsync(); } catch { await _dbContext.Database.RollbackTransactionAsync(); throw; } } }
2. microservice_A:订阅事件并创建Account
// 定义Account创建完成事件 public record AccountCreated(Guid AccountId, string AccountName, Guid CustomerId); // 事件消费者 public class CustomerCreatedPendingAccountConsumer : IConsumer<CustomerCreatedPendingAccount> { private readonly AppDbContext _dbContext; private readonly IPublishEndpoint _publishEndpoint; public CustomerCreatedPendingAccountConsumer(AppDbContext dbContext, IPublishEndpoint publishEndpoint) { _dbContext = dbContext; _publishEndpoint = publishEndpoint; } public async Task Consume(ConsumeContext<CustomerCreatedPendingAccount> context) { var message = context.Message; // 用聚合工厂创建Account(符合microservice_A业务规则) var account = AccountFactory.CreateAccount(message.CustomerName); await _dbContext.Accounts.AddAsync(account); await _dbContext.SaveChangesAsync(); // 发布Account创建完成事件 await _publishEndpoint.Publish(new AccountCreated(account.AccountId, account.AccountName, message.CustomerId)); } }
3. microservice_B:订阅事件更新Customer关联
public class AccountCreatedConsumer : IConsumer<AccountCreated> { private readonly AppDbContext _dbContext; public AccountCreatedConsumer(AppDbContext dbContext) { _dbContext = dbContext; } public async Task Consume(ConsumeContext<AccountCreated> context) { var message = context.Message; var customer = await _dbContext.Customers.FirstOrDefaultAsync(c => c.CustomerId == message.CustomerId); if (customer == null) { // 记录日志并转入死信队列,等待重试或人工处理 throw new InvalidOperationException($"Customer {message.CustomerId} not found"); } // 更新Customer的AccountId(符合microservice_B业务规则) customer.AssignAccount(message.AccountId); await _dbContext.SaveChangesAsync(); } }
4. MassTransit配置(RabbitMQ为例)
在Program.cs中配置消息队列、重试和死信:
builder.Services.AddMassTransit(x => { // 注册microservice_B的消费者 x.AddConsumer<AccountCreatedConsumer>(); x.UsingRabbitMq((context, cfg) => { cfg.Host("rabbitmq://localhost"); cfg.ReceiveEndpoint("customer-service-queue", e => { e.ConfigureConsumer<AccountCreatedConsumer>(context); // 配置3次间隔重试 e.UseMessageRetry(r => r.Interval(3, TimeSpan.FromSeconds(5))); // 配置死信队列 e.DeadLetterExchange = "dead-letter-exchange"; e.DeadLetterRoutingKey = "dead-letter"; }); }); });
方案二:同步通信+补偿机制(仅适用于强一致性要求极高的场景)
若必须用同步调用,需通过TCC事务模式(Try-Confirm-Cancel)实现补偿,保证跨服务操作一致性。
具体流程
- microservice_B调用microservice_A的
ReserveAccount接口,创建"待确认"状态的Account - microservice_B创建关联该AccountId的CUSTOMER记录
- 若CUSTOMER创建成功,调用microservice_A的
ConfirmAccount接口,将Account状态改为"正常" - 若CUSTOMER创建失败,调用microservice_A的
CancelAccount接口,删除或标记Account为无效
.NET技术栈实现示例
使用Polly实现重试/熔断,结合自定义补偿逻辑:
public class CustomerService { private readonly AppDbContext _dbContext; private readonly HttpClient _httpClient; private readonly IAsyncPolicy<HttpResponseMessage> _retryPolicy; public CustomerService(AppDbContext dbContext, HttpClient httpClient) { _dbContext = dbContext; _httpClient = httpClient; // 配置Polly指数退避重试策略 _retryPolicy = Policy .HandleResult<HttpResponseMessage>(r => !r.IsSuccessStatusCode) .WaitAndRetryAsync(3, retryAttempt => TimeSpan.FromSeconds(Math.Pow(2, retryAttempt))); } public async Task CreateCustomerAsync(string customerName) { Guid accountId = Guid.Empty; bool accountConfirmed = false; try { // 1. 调用microservice_A预留Account var reserveResponse = await _retryPolicy.ExecuteAsync(() => _httpClient.PostAsJsonAsync("http://microservice-a/api/accounts/reserve", new { AccountName = customerName })); reserveResponse.EnsureSuccessStatusCode(); accountId = await reserveResponse.Content.ReadFromJsonAsync<Guid>(); // 2. 创建关联Account的Customer var customer = CustomerFactory.CreateCustomer(customerName, accountId); await _dbContext.Customers.AddAsync(customer); await _dbContext.SaveChangesAsync(); // 3. 确认Account生效 var confirmResponse = await _retryPolicy.ExecuteAsync(() => _httpClient.PostAsJsonAsync($"http://microservice-a/api/accounts/{accountId}/confirm", new {})); confirmResponse.EnsureSuccessStatusCode(); accountConfirmed = true; } catch { // 4. 补偿:若Account已预留未确认,调用取消接口 if (!accountConfirmed && accountId != Guid.Empty) { await _httpClient.PostAsJsonAsync($"http://microservice-a/api/accounts/{accountId}/cancel", new {}); } throw; } } }
最佳实践总结
- 优先选择事件驱动:符合微服务松耦合原则,避免同步调用的依赖和超时风险
- 本地事务+事件原子性:保证本地数据变更与事件发布同时成功或失败
- 重试+死信兜底:处理临时故障,避免业务流程卡在中间状态
- 状态显性化:为聚合根增加状态字段(如"待关联"、"正常"),清晰反映跨服务操作进度
- 避免无效null值:用特定标记代替null,严格遵循业务规则约束
内容的提问来源于stack exchange,提问作者Ali Adlavaran
相关产品推荐
相关产品推荐

