如何在MassTransit的RoutingSlipCompleted消费方法中获取Activity数据
如何将Activity中的数据传递到RoutingSlipCompleted消费方法中
我在Activity中存储了数据,希望能在RoutingSlipCompleted消费方法中获取这些数据。我知道可以通过CompletedWithVariables在Activity之间传递数据,但不清楚如何将Activity中的数据传递到RoutingSlipCompleted中。
CheckInventoryActivity 实现
public class CheckInventoryActivity : IActivity<ICheckInventoryRequest, CheckInventoryRequestCompensate> { private readonly IInventoryService _inventoryService; private readonly IEndpointNameFormatter _formatter; public CheckInventoryActivity(IInventoryService inventoryService, IEndpointNameFormatter formatter) { _inventoryService = inventoryService; _formatter = formatter; } public async Task<ExecutionResult> Execute(ExecuteContext<ICheckInventoryRequest> context) { CheckInventoryRequest model = new CheckInventoryRequest() { PartCode = context.Arguments.PartCode }; var response = await _inventoryService.CheckInventory(model); var checkInventoryResponse = new CheckInventoryResponse() { PartCode = response.Data.PartCode ?? model.PartCode, Id = response.Data.Id ?? 0, InventoryCount = response.Data.InventoryCount ?? 0 }; var checkInventoryCustomActionResult = new CustomActionResult<CheckInventoryResponse>() { Data = checkInventoryResponse, IsSuccess = true, ResponseDesc = "success", ResponseType = 0 }; var result = response.IsSuccess; if (!result) return context.CompletedWithVariables<CheckInventoryRequestCompensate>( new { Result = result, LogDate = DateTime.Now, MethodName = "CheckInventoryActivity", }, new { Result = result, LogDate = DateTime.Now, MethodName = "CheckInventoryActivity", CheckInventoryCustomActionResult = checkInventoryCustomActionResult }); var queueName = _formatter.ExecuteActivity<CallSuccessActivity, ISuccessRequest>(); var uri = new Uri($"queue:{queueName}"); return context.ReviseItinerary(x => x.AddActivity("CallSuccessActivity", uri, new { LogDate = DateTime.Now, MethodName = "CheckInventoryActivity", CheckInventoryCustomActionResult = checkInventoryCustomActionResult })); } }
在CallSuccessActivity中获取数据的方式
return context.ReviseItinerary(x => x.AddActivity("CallSuccessActivity", uri, new { LogDate = DateTime.Now, MethodName = "CheckInventoryActivity", CheckInventoryCustomActionResult = checkInventoryCustomActionResult }));
CallSuccessActivity 实现
public class CallSuccessActivity : IExecuteActivity<ISuccessRequest> { private readonly IRequestClient<ISuccessRequest> _requestClient; public CallSuccessActivity(IRequestClient<ISuccessRequest> requestClient) { _requestClient = requestClient; } public async Task<ExecutionResult> Execute(ExecuteContext<ISuccessRequest> context) { var iModel = context.Arguments; var model = new SuccessRequest() { LogDate = iModel.LogDate, MethodName = iModel.MethodName, CheckInventoryCustomActionResult = iModel.CheckInventoryCustomActionResult }; //CustomActionResult< CheckInventoryResponse > CheckInventoryResponse = new (); var rabbitResult = await _requestClient.GetResponse<CustomActionResult<SuccessResponse>>(model); return context.Completed(); } }
当前RoutingSlipCompleted消费方法实现
public async Task Consume(ConsumeContext<RoutingSlipCompleted> context) { var requestId = context.Message.GetVariable<Guid?>(nameof(ConsumeContext.RequestId)); var checkInventoryResponseModel = context.Message.Variables["CheckInventoryResponse"]; var responseAddress = context.Message.GetVariable<Uri>(nameof(ConsumeContext.ResponseAddress)); var request = context.Message.GetVariable<ICheckInventoryRequest>("Model"); throw new NotImplementedException(); }
解决方案
要把Activity中的数据传递到RoutingSlipCompleted,核心是将数据存入路由单的全局变量,路由单完成时会携带这些全局变量,直接从RoutingSlipCompleted的Variables集合中读取即可。
步骤1:在Activity中存入全局变量
根据Activity的不同返回场景,选择对应的方式添加全局变量:
场景1:Activity直接完成(使用CompletedWithVariables)
修改CheckInventoryActivity中返回CompletedWithVariables的代码,确保全局变量包含目标数据:
return context.CompletedWithVariables<CheckInventoryRequestCompensate>( // 补偿逻辑用的变量(可选) new { Result = result, LogDate = DateTime.Now, MethodName = "CheckInventoryActivity", }, // 全局变量:会被携带到RoutingSlipCompleted new Dictionary<string, object> { ["Result"] = result, ["LogDate"] = DateTime.Now, ["MethodName"] = "CheckInventoryActivity", ["CheckInventoryCustomActionResult"] = checkInventoryCustomActionResult });
场景2:修改路由单(使用ReviseItinerary)
在添加新Activity的同时,通过SetVariable方法将数据存入全局变量:
return context.ReviseItinerary(x => x.AddActivity("CallSuccessActivity", uri, new { LogDate = DateTime.Now, MethodName = "CheckInventoryActivity", CheckInventoryCustomActionResult = checkInventoryCustomActionResult }) // 添加全局变量 .SetVariable("CheckInventoryCustomActionResult", checkInventoryCustomActionResult));
场景3:在CallSuccessActivity中传递数据
如果需要在CallSuccessActivity完成后传递数据,同样使用CompletedWithVariables存入全局变量:
return context.CompletedWithVariables(new { CheckInventoryCustomActionResult = iModel.CheckInventoryCustomActionResult });
步骤2:在RoutingSlipCompleted中读取数据
修改Consume方法,从context.Message.Variables中取出目标数据:
public async Task Consume(ConsumeContext<RoutingSlipCompleted> context) { var requestId = context.Message.GetVariable<Guid?>(nameof(ConsumeContext.RequestId)); // 读取全局变量中的CheckInventoryCustomActionResult if (context.Message.Variables.TryGetValue("CheckInventoryCustomActionResult", out var value)) { var checkInventoryResult = value as CustomActionResult<CheckInventoryResponse>; // 这里编写处理该数据的业务逻辑 } var responseAddress = context.Message.GetVariable<Uri>(nameof(ConsumeContext.ResponseAddress)); var request = context.Message.GetVariable<ICheckInventoryRequest>("Model"); // 替换NotImplementedException为实际业务代码 }
关键注意事项
- 全局变量会跟随整个路由流程,最终完整出现在
RoutingSlipCompleted消息中。 - 复杂对象需要支持序列化(比如添加
[Serializable]属性,或确保符合MassTransit序列化器的要求)。 - 无论是
SetVariable还是CompletedWithVariables的全局参数,都能将数据存入路由单的全局变量集合。
内容的提问来源于stack exchange,提问作者Ali Eshghi
相关产品推荐
相关产品推荐

