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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 21:40:19