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

使用RoutingSlipProxy的MassTransit请求客户端故障排查与示例求助

ASP.NET Core 结合 RequestClient 与 RoutingSlipRequestProxy/RoutingSlipResponseProxy 示例及问题解决

一、完整配置与代码示例

1. ASP.NET Core 依赖注入配置

在 Program.cs 中完成 MassTransit 核心配置,包含路由 slip 代理、请求客户端、消息调度器(解决 Saga 超时问题)的注册:

using MassTransit;
using MassTransit.RoutingSlips;

var builder = WebApplication.CreateBuilder(args);

// 注册 MassTransit 核心服务
builder.Services.AddMassTransit(x =>
{
    // 注册响应消费者
    x.AddConsumer<CreateESimOrderCommandResponseConsumer>();

    // 注册 RoutingSlip 业务活动(替换为你的实际活动)
    x.AddActivity<ProcessESimOrderActivity, ProcessESimOrderActivityArguments>();

    // 启用路由 slip 请求/响应代理
    x.AddRoutingSlipRequestProxy();
    x.AddRoutingSlipResponseProxy();

    // 注册消息调度器,处理请求超时逻辑
    x.AddMessageScheduler(new Uri("queue:scheduler"));

    x.UsingRabbitMq((context, cfg) =>
    {
        cfg.Host("localhost", "/", h =>
        {
            h.Username("guest");
            h.Password("guest");
        });

        // 自动配置所有端点(包括活动、消费者、调度器)
        cfg.ConfigureEndpoints(context);
        // 全局启用消息调度器
        cfg.UseMessageScheduler(new Uri("queue:scheduler"));

        // 手动绑定响应消费者到指定队列(可选,若自动配置未生效)
        cfg.ReceiveEndpoint("create-esim-order-response", e =>
        {
            e.ConfigureConsumer<CreateESimOrderCommandResponseConsumer>(context);
        });
    });
});

// 注册请求客户端,供控制器/Saga 调用
builder.Services.AddRequestClient<CreateESimOrderCommand>();

var app = builder.Build();
app.Run();

2. 消息与活动参数定义

// 请求消息,需实现 CorrelatedBy<Guid> 确保路由 slip 与请求关联
public record CreateESimOrderCommand : CorrelatedBy<Guid>
{
    public Guid CorrelationId { get; init; }
    public string Iccid { get; init; } // 业务字段
}

// 响应消息
public record CreateESimOrderCommandResponse
{
    public Guid CorrelationId { get; init; }
    public bool Success { get; init; }
    public string Message { get; init; }
}

// RoutingSlip 活动执行参数
public record ProcessESimOrderActivityArguments
{
    public string Iccid { get; init; }
}

3. RoutingSlip 业务活动实现

public class ProcessESimOrderActivity : IExecuteActivity<ProcessESimOrderActivityArguments>
{
    public async Task<ExecutionResult> Execute(ExecuteContext<ProcessESimOrderActivityArguments> context)
    {
        // 执行业务逻辑:比如调用第三方接口创建 eSIM 订单
        var iccid = context.Arguments.Iccid;
        
        // 模拟业务处理延迟
        await Task.Delay(1000);

        // 返回活动执行结果,可被后续活动或响应代理使用
        return context.Completed(new { Iccid = iccid, Success = true });
    }
}

4. 响应消费者实现

public class CreateESimOrderCommandResponseConsumer : IConsumer<CreateESimOrderCommandResponse>
{
    public async Task Consume(ConsumeContext<CreateESimOrderCommandResponse> context)
    {
        // 处理响应逻辑:比如更新订单状态、通知前端等
        Console.WriteLine($"收到 eSIM 订单响应: 状态={context.Message.Success}, 消息={context.Message.Message}");
        await Task.CompletedTask;
    }
}

5. 控制器中使用 RequestClient 调用路由 slip

[ApiController]
[Route("api/esim")]
public class ESimController : ControllerBase
{
    private readonly IRequestClient<CreateESimOrderCommand> _requestClient;

    public ESimController(IRequestClient<CreateESimOrderCommand> requestClient)
    {
        _requestClient = requestClient;
    }

    [HttpPost("order")]
    public async Task<IActionResult> CreateOrder(string iccid)
    {
        var correlationId = Guid.NewGuid();
        
        // 构建路由 slip,指定要执行的活动
        var routingSlip = new RoutingSlipBuilder(correlationId)
            .AddActivity("ProcessESimOrder", new Uri("queue:process-esim-order"), new ProcessESimOrderActivityArguments { Iccid = iccid })
            .Build();

        // 发送请求并等待响应
        var response = await _requestClient.GetResponse<CreateESimOrderCommandResponse>(routingSlip);

        return Ok(response.Message);
    }
}

6. Saga 中调用路由 slip 请求代理

public class ESimOrderSaga :
    ISaga,
    InitiatedBy<InitiateESimOrderSagaCommand>,
    Orchestrates<CreateESimOrderCommandResponse>
{
    public Guid CorrelationId { get; set; }
    public bool OrderCreated { get; set; }

    public async Task Consume(ConsumeContext<InitiateESimOrderSagaCommand> context)
    {
        // 构建路由 slip,使用 Saga 的 CorrelationId 确保关联
        var routingSlip = new RoutingSlipBuilder(CorrelationId)
            .AddActivity("ProcessESimOrder", new Uri("queue:process-esim-order"), new ProcessESimOrderActivityArguments { Iccid = context.Message.Iccid })
            .Build();

        // 从上下文获取请求客户端,自动继承调度器配置
        var requestClient = context.GetRequestClient<CreateESimOrderCommand>();
        // 指定超时时间,依赖消息调度器实现
        await requestClient.GetResponse<CreateESimOrderCommandResponse>(routingSlip, context.CancellationToken, TimeSpan.FromSeconds(30));
    }

    public async Task Consume(ConsumeContext<CreateESimOrderCommandResponse> context)
    {
        OrderCreated = context.Message.Success;
        await Task.CompletedTask;
    }
}

二、问题排查与解决

1. 请求进入 SKIP 队列,响应消费者未执行

  • 活动未正确注册/绑定:确保使用 AddActivity 注册业务活动,且 ConfigureEndpoints 自动创建了活动对应的队列;若自动配置失效,可手动添加 ReceiveEndpoint 绑定活动队列。
  • 请求与路由 slip 关联失效:请求消息必须实现 CorrelatedBy<Guid>,且路由 slip 的 CorrelationId 与请求的 CorrelationId 完全一致,否则响应代理无法将路由 slip 完成消息转换为请求响应。
  • 响应消费者未正确绑定:检查 AddConsumer 和 ReceiveEndpoint 配置,确保响应消息的队列与消费者关联,且消费者类未遗漏 IConsumer<T> 接口实现。
  • 活动执行失败:查看 MassTransit 日志,若活动抛出未捕获异常,消息会进入错误队列而非触发响应,需修复活动业务逻辑中的异常。

2. Saga 中调用时抛出"A request timeout was specified but no message scheduler was specified or available"

  • 未注册消息调度器:在 MassTransit 配置中必须添加 AddMessageScheduler 和 UseMessageScheduler,指定调度器队列(如示例中的 queue:scheduler),超时逻辑依赖调度器发送延迟消息。
  • 调度器端点未启动:确保 ConfigureEndpoints 自动创建了调度器端点,若手动配置端点,需确保调度器队列被正确初始化。
  • Saga 中错误获取 RequestClient:避免在 Saga 构造函数中注入 IRequestClient,应通过 context.GetRequestClient<T>() 获取,确保使用当前上下文的调度器配置。

内容的提问来源于stack exchange,提问作者Robert Mircea

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 05:16:23