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

如何在MassTransit的IConsumer中注入DurableTaskClient?

问题描述

我有一个Azure Function,使用Orchestrator从外部源同步产品。此前该函数通过TimerTrigger或HttpTrigger传入InventoryId触发,仅能同步指定产品。

现在需要为其添加队列支持,但无法在IConsumer中注入DurableTaskClient。目前实现的代码可以运行,但完全未用到MassTransit Consumer。由于生产者端全程使用MassTransit,我们不想修改现有行为,也不想发送不带完整消息体的普通消息。

想请教两个问题:

  • 是否有办法在MassTransit Consumer中注入DurableTaskClient?
  • 若无法注入,有没有比当前更简洁的消息反序列化方式?

当前代码如下:

[Function(nameof(MessageReceiverHandler))]
public async Task Run(
    [ServiceBusTrigger("%QueueName%", Connection = "ServiceBusConnection")] ServiceBusReceivedMessage message, [DurableClient] DurableTaskClient client, CancellationToken cancellationToken)
{
    var str = message.Body.ToString();

    var definition = new { message = new { inventoryId =3} };
    var res = JsonConvert.DeserializeAnonymousType(str, definition);

    await receiver.HandleConsumer<ProductSyncConsumer>(options.Value.QueueName, message, cancellationToken);

    string instanceId = await client.ScheduleNewOrchestrationInstanceAsync(nameof(DurablePartnerSyncActions), res.message.inventoryId, cancellationToken);

    logger.LogInformation("Started Individual orchestration with ID = '{instanceId}'.", instanceId);
   
}

public class ProductSyncConsumer : IConsumer<ExternalProductSyncSingle>
{
    private readonly ILogger<ProductSyncConsumer> _logger;

    public ProductSyncConsumer(ILogger<ProductSyncConsumer> logger)
    {
      
        _logger = logger;
    }

    public  Task Consume(ConsumeContext<ExternalProductSyncSingle> context)
    {
        return Task.FromResult(context.Message.InventoryId);
       // DurableTaskClient client = default(DurableTaskClient);
       // string instanceId = await client.ScheduleNewOrchestrationInstanceAsync(nameof(DurablePartnerSyncActions), context.Message.InventoryId);

        _logger.LogInformation("Started Individual orchestration with ID = '{instanceId}'.", 5);
      //  return await client.CreateCheckStatusResponseAsync(req, instanceId);
    }
}
解决方案

一、在MassTransit Consumer中注入DurableTaskClient

MassTransit在Azure Functions中支持完整的依赖注入,DurableTaskClient可直接注入到Consumer构造函数中(Azure Functions会自动注册该服务到DI容器)。

修改Consumer代码

public class ProductSyncConsumer : IConsumer<ExternalProductSyncSingle>
{
    private readonly ILogger<ProductSyncConsumer> _logger;
    private readonly DurableTaskClient _durableClient;

    // 直接注入DurableTaskClient
    public ProductSyncConsumer(ILogger<ProductSyncConsumer> logger, DurableTaskClient durableClient)
    {
        _logger = logger;
        _durableClient = durableClient;
    }

    public async Task Consume(ConsumeContext<ExternalProductSyncSingle> context)
    {
        string instanceId = await _durableClient.ScheduleNewOrchestrationInstanceAsync(
            nameof(DurablePartnerSyncActions), 
            context.Message.InventoryId, 
            context.CancellationToken);

        _logger.LogInformation("Started Individual orchestration with ID = '{instanceId}'.", instanceId);
    }
}

配置MassTransit自动路由

推荐使用MassTransit的Azure Functions集成,让框架自动处理消息路由,无需手动编写触发器函数。在Program.cs中添加配置:

builder.Services.AddMassTransit(x =>
{
    x.AddConsumer<ProductSyncConsumer>();

    x.UsingAzureServiceBus((context, cfg) =>
    {
        cfg.Host(builder.Configuration["ServiceBusConnection"]);

        cfg.ReceiveEndpoint(builder.Configuration["QueueName"], e =>
        {
            e.ConfigureConsumer<ProductSyncConsumer>(context);
        });
    });
});

配置完成后,MassTransit会自动创建对应的ServiceBus触发器,将消息路由到ProductSyncConsumer,同时DI容器会自动注入DurableTaskClient。

二、更简洁的消息反序列化方式

若暂时不想调整MassTransit配置,可直接将消息反序列化为目标类型,替代匿名类型的冗余写法:

直接反序列化为目标DTO

[Function(nameof(MessageReceiverHandler))]
public async Task Run(
    [ServiceBusTrigger("%QueueName%", Connection = "ServiceBusConnection")] ServiceBusReceivedMessage message, 
    [DurableClient] DurableTaskClient client, 
    CancellationToken cancellationToken)
{
    // 直接反序列化为ExternalProductSyncSingle类型
    var syncMessage = JsonConvert.DeserializeObject<ExternalProductSyncSingle>(message.Body.ToString());

    string instanceId = await client.ScheduleNewOrchestrationInstanceAsync(
        nameof(DurablePartnerSyncActions), 
        syncMessage.InventoryId, 
        cancellationToken);

    logger.LogInformation("Started Individual orchestration with ID = '{instanceId}'.", instanceId);
}

处理嵌套结构的消息

若消息外层有嵌套包装(如{ "message": {...} }),可定义对应的包装类:

public class SyncMessageWrapper
{
    public ExternalProductSyncSingle Message { get; set; }
}

// 反序列化时使用包装类
var wrapper = JsonConvert.DeserializeObject<SyncMessageWrapper>(message.Body.ToString());
var inventoryId = wrapper.Message.InventoryId;

这种方式比匿名类型更清晰,也便于后续维护。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 16:12:49