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

基于最终一致性的一对一关系:CreateCustomer用例实现问询

微服务最终一致性实践:CreateCustomer用例实现方案

问题1:CreateCustomer用例应由哪个微服务负责实现?

你的判断完全正确,microservice_B应该主导这个用例。因为CUSTOMER是microservice_B的核心聚合根,CreateCustomer是围绕CUSTOMER生命周期展开的业务操作,由归属的微服务负责更符合领域驱动设计(DDD)的聚合根职责边界原则。

问题2:CreateCustomer用例如何实现?有哪些最佳实践?

结合最终一致性原则和.NET技术栈生态,针对你提出的两个方案,以下是优化后的落地实践:

方案一:事件驱动的最终一致性(推荐)

异步事件驱动是实现跨微服务最终一致性的标准方案,需补充状态跟踪和重试/补偿机制解决中间状态卡顿问题,避免业务流程中断。

具体流程

  1. microservice_B接收CreateCustomer请求,创建**状态为"待关联账户"**的CUSTOMER记录(用0/Guid.Empty等标记代替null,避免违反"每个CUSTOMER必须关联ACCOUNT"的业务规则,仅标识关联未完成)
  2. microservice_B在本地事务内发布CustomerCreatedPendingAccount事件(包含CustomerId、CustomerName)
  3. microservice_A监听事件,通过自身聚合工厂创建符合业务规则的ACCOUNT,生成AccountId
  4. microservice_A发布AccountCreated事件(包含AccountId、AccountName、关联的CustomerId)
  5. microservice_B监听事件,更新对应CUSTOMER的AccountId,并将状态改为"正常"
  6. 配置死信队列和定时补偿任务:若某一步事件处理失败,将事件转入死信队列,通过定时任务重试;若重试仍失败,触发人工介入排查。

.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)实现补偿,保证跨服务操作一致性。

具体流程

  1. microservice_B调用microservice_A的ReserveAccount接口,创建"待确认"状态的Account
  2. microservice_B创建关联该AccountId的CUSTOMER记录
  3. 若CUSTOMER创建成功,调用microservice_A的ConfirmAccount接口,将Account状态改为"正常"
  4. 若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;
        }
    }
}

最佳实践总结

  1. 优先选择事件驱动:符合微服务松耦合原则,避免同步调用的依赖和超时风险
  2. 本地事务+事件原子性:保证本地数据变更与事件发布同时成功或失败
  3. 重试+死信兜底:处理临时故障,避免业务流程卡在中间状态
  4. 状态显性化:为聚合根增加状态字段(如"待关联"、"正常"),清晰反映跨服务操作进度
  5. 避免无效null值:用特定标记代替null,严格遵循业务规则约束

内容的提问来源于stack exchange,提问作者Ali Adlavaran

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 18:35:02