在Saga模式中捕获并处理Kafka消费者抛出的异常
问题原因
你编写的OrderManagementSystemConsumer是独立于Saga的消息消费者,它和Saga的事件处理逻辑属于两个并行的消费链路。当OrderRequestEvent消息到达Kafka主题时,这个独立消费者和Saga的事件消费者会分别处理该消息,因此独立消费者抛出的异常只会影响自身的消费重试逻辑,不会中断或影响Saga的流程执行。
解决方案
有两种常见的处理方式,可根据架构需求选择:
方案一:将验证逻辑整合到Saga流程中
直接把验证逻辑移到Saga的OrderRequestedEvent处理步骤里,利用Saga的Catch机制捕获异常,执行状态修改、重试消息发送等操作:
Initially( When(OrderRequestedEvent) .Then(context => LogContext.Info?.Log("Initializing saga: {0}", context.Saga.CorrelationId)) .InitializeSaga() // 将验证逻辑嵌入Saga流程 .ThenAsync(async context => { var apiService = context.GetService<IApiService>(); if (await apiService.ValidateIncomingRequestAsync(context.Message)) throw new ArgumentException("Something wrong just happened"); }) // 捕获验证异常并执行自定义逻辑 .Catch<ArgumentException>(ex => ex .Then(context => LogContext.Info?.Log("Validation failed for saga: {0}, error: {1}", context.Saga.CorrelationId, ex.Message)) .Then(context => { // 修改Saga状态为验证失败 context.Saga.Status = SagaStatus.ValidationFailed; }) // 发送重试消息到指定主题 .Send(context => new RetryOrderValidationEvent { CorrelationId = context.Saga.CorrelationId, OrderRequest = context.Message }) // 切换到验证失败状态,后续可处理重试逻辑 .TransitionTo(ValidationFailed)) // 验证成功后的正常流程 .Then(context => LogContext.Info?.Log("Validating Customer: {0}", context.Saga.CorrelationId)) .SendingToCustomerValidation().LogSaga() .TransitionTo(ValidatingCustomer) );
这种方式的优势是逻辑集中在Saga内部,无需额外独立消费者,异常处理和流程控制更统一。
方案二:保留独立消费者,通过事件通知Saga处理异常
如果必须保留独立的验证消费者,可以在验证失败时发送“验证失败”事件,让Saga监听该事件并执行后续逻辑:
修改独立消费者代码
public async Task Consume(ConsumeContext<OrderRequestEvent> context) { ArgumentNullException.ThrowIfNull(context, nameof(context)); try { if (await this.apiService.ValidateIncomingRequestAsync(context.Message)) throw new ArgumentException("Something wrong just happened"); // 验证成功,发送事件通知Saga继续流程 await context.Publish(new OrderValidationSuccessEvent { CorrelationId = context.Message.CorrelationId }); } catch (ArgumentException ex) { // 验证失败,发送事件通知Saga处理异常 await context.Publish(new OrderValidationFailedEvent { CorrelationId = context.Message.CorrelationId, ErrorMessage = ex.Message, OrderRequest = context.Message }); } }
修改Saga代码,添加失败事件处理
// 初始状态:等待验证结果 Initially( When(OrderRequestedEvent) .Then(context => LogContext.Info?.Log("Initializing saga: {0}", context.Saga.CorrelationId)) .InitializeSaga() .TransitionTo(AwaitingValidationResult)); // 处理验证结果的状态分支 During(AwaitingValidationResult, // 验证成功,执行正常流程 When(OrderValidationSuccessEvent) .Then(context => LogContext.Info?.Log("Validating Customer: {0}", context.Saga.CorrelationId)) .SendingToCustomerValidation().LogSaga() .TransitionTo(ValidatingCustomer), // 验证失败,执行异常处理逻辑 When(OrderValidationFailedEvent) .Then(context => { context.Saga.Status = SagaStatus.ValidationFailed; LogContext.Info?.Log("Validation failed for saga: {0}, error: {1}", context.Saga.CorrelationId, context.Message.ErrorMessage); }) .Send(context => new RetryOrderValidationEvent { CorrelationId = context.Saga.CorrelationId, OrderRequest = context.Message.OrderRequest }) .TransitionTo(ValidationFailed) );
这种方式适合验证逻辑需要独立于Saga的场景,通过事件驱动实现Saga和验证逻辑的解耦。
内容的提问来源于stack exchange,提问作者Stanley_Marsh
相关产品推荐
相关产品推荐

