如何在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
相关产品推荐
相关产品推荐

