MassTransit Saga技术问询:状态机内调用API并触发多事件
嘿,作为经常用MassTransit和Automatonymous开发Saga的开发者,我来帮你搞定这个问题!你已经理清了响应收集的逻辑,现在就差在状态机里调用API获取县列表并触发后续事件这一步,其实实现起来很直观,我给你拆解思路和代码示例:
核心思路概述
在你的Saga状态机中,初始状态收到触发事件后,直接集成API调用获取目标州的县列表,遍历列表为每个县发布计算事件,之后进入等待所有计算完成的状态。整个流程可以直接在状态机的ThenAsync环节完成,不需要额外的复杂组件。
具体实现步骤
- 定义必要的事件、Saga数据模型,用来承载状态和计算数据
- 在状态机的初始转换逻辑中,注入API客户端并调用接口获取县列表
- 遍历县列表,为每个县发布计算事件,并将待处理县记录到Saga数据中
- 后续处理每个县的计算完成事件,收集结果并判断是否所有任务完成
- 所有计算完成后生成最终统计报告
关键代码示例
1. 定义事件与Saga数据模型
// 触发Saga启动的初始事件 public record InitiateStateFoodConsumptionCalculation(Guid CorrelationId, string StateCode); // 为单个县发布的计算事件 public record CalculateCountyFoodConsumption(Guid CorrelationId, string CountyCode); // 单个县计算完成后的回调事件 public record CountyFoodConsumptionCalculated(Guid CorrelationId, string CountyCode, decimal TotalConsumption); // Saga数据类,保存状态、待处理县列表和计算结果 public class FoodConsumptionSagaData : SagaStateMachineInstance { public Guid CorrelationId { get; set; } public string CurrentState { get; set; } public string TargetStateCode { get; set; } public List<string> PendingCountyCodes { get; set; } = new(); public Dictionary<string, decimal> CountyConsumptionResults { get; set; } = new(); }
2. 实现状态机逻辑
// 假设你有一个已封装好的地理州API客户端 public interface IGeoStateApiClient { Task<List<CountyDto>> GetCountiesByStateCode(string stateCode); } public class FoodConsumptionSaga : MassTransitStateMachine<FoodConsumptionSagaData> { // 注入API客户端,MassTransit支持状态机构造函数注入依赖 public FoodConsumptionSaga(IGeoStateApiClient geoStateApiClient) { InstanceState(x => x.CurrentState); // 定义事件关联规则 Event(() => InitiateCalculation, x => x.CorrelateById(m => m.Message.CorrelationId)); Event(() => CountyCalculationCompleted, x => x.CorrelateById(m => m.Message.CorrelationId)); // 初始状态处理:收到启动事件后调用API、发布县计算事件 Initially( When(InitiateCalculation) .ThenAsync(async context => { // 调用API获取目标州的所有县 var counties = await geoStateApiClient.GetCountiesByStateCode(context.Message.StateCode); // 保存状态到Saga数据 context.Instance.TargetStateCode = context.Message.StateCode; context.Instance.PendingCountyCodes.AddRange(counties.Select(c => c.Code)); // 为每个县发布计算事件 foreach (var county in counties) { await context.Publish(new CalculateCountyFoodConsumption( context.Instance.CorrelationId, county.Code)); } }) .TransitionTo(WaitingForCountyCalculations) ); // 等待所有县计算完成的状态 During(WaitingForCountyCalculations, When(CountyCalculationCompleted) .Then(context => { // 记录当前县的计算结果 context.Instance.CountyConsumptionResults[context.Message.CountyCode] = context.Message.TotalConsumption; // 从待处理列表移除已完成的县 context.Instance.PendingCountyCodes.Remove(context.Message.CountyCode); }) // 判断是否所有县都处理完成 .If(context => !context.Instance.PendingCountyCodes.Any(), x => x .ThenAsync(async context => { // 生成最终统计报告 var finalReport = new StateFoodConsumptionReport { StateCode = context.Instance.TargetStateCode, TotalStateConsumption = context.Instance.CountyConsumptionResults.Values.Sum(), CountyBreakdown = context.Instance.CountyConsumptionResults }; // 发布报告完成事件,或者直接写入数据库/通知外部系统 await context.Publish(finalReport); }) .TransitionTo(Completed)) ); SetCompletedWhenFinalized(); } // 定义状态节点 public State WaitingForCountyCalculations { get; private set; } public State Completed { get; private set; } // 定义事件 public Event<InitiateStateFoodConsumptionCalculation> InitiateCalculation { get; private set; } public Event<CountyFoodConsumptionCalculated> CountyCalculationCompleted { get; private set; } }
补充注意事项
- 依赖注入:确保你的
IGeoStateApiClient已注册到DI容器中,MassTransit会自动将其注入到状态机构造函数。 - 错误处理:可以给API调用和事件发布添加故障处理逻辑,比如用
OnFaulted捕获API调用失败的情况,触发重试或补偿流程。 - 持久化配置:记得为Saga数据配置持久化(比如Entity Framework、MongoDB),确保Saga状态和数据在服务重启后不会丢失。
内容的提问来源于stack exchange,提问作者Ceemah Four
相关产品推荐
相关产品推荐

