如何通过Azure Functions输出绑定与IAsyncCollector设置Service Bus消息的MessageId
如何通过Azure Functions输出绑定与IAsyncCollector设置Service Bus消息的MessageId?
你当前代码里用的IAsyncCollector<string>只能发送纯字符串内容,没法直接设置MessageId这类消息属性。要实现这个需求,你需要把泛型类型换成Service Bus的消息对象类型,具体分两种情况(取决于你用的Service Bus扩展包版本):
情况1:使用老版Service Bus扩展(v5及以下,基于Microsoft.Azure.ServiceBus SDK)
- 把函数参数里的
IAsyncCollector<string>替换为IAsyncCollector<Message> - 构造
Message实例时,直接设置MessageId属性,再把序列化后的内容作为消息体
修改后的代码示例:
[FunctionName(nameof(ReceiveTransaction))] [OpenApiOperation(nameof(ReceiveTransaction), "Transaction", Description = "Receives a transaction and adds to the Transaction Platform")] [OpenApiRequestBody("application/json", typeof(string), Description = "Parameters", Required = true)] [OpenApiResponseWithoutBody(HttpStatusCode.BadRequest, Description = "Request was invalid and could not be processed")] [OpenApiResponseWithoutBody(HttpStatusCode.OK, Description = "Request was valid - transaction posted to platform")] public async Task<IActionResult> ReceiveTransaction( [HttpTrigger(AuthorizationLevel.Anonymous, "post", Route = "v1/receiveTransaction")] HttpRequest req, // 替换为IAsyncCollector<Message> [ServiceBus("%Transactions:SB:Transactions:IngressQueue%", Connection = "Transactions:SB:Transactions")] IAsyncCollector<Message> collector, ILogger logger) { try { using var sr = new StreamReader(req.Body); var rawJson = await sr.ReadToEndAsync(); var tran = JsonConvert.DeserializeObject<RetailTransaction>(rawJson, JsonSerialisationUtils.SerialiserSettings); tran.UpdateExtendedProperties(); tran.SetPlatformReceived(); logger.LogInformation(StandardEvents.SuccessEvent, "Incoming transaction: {transactionNumber} validated", tran.TransactionNumber); // 构造Message对象并设置MessageId var messageBody = JsonConvert.SerializeObject(tran, JsonSerialisationUtils.SerialiserSettings); var serviceBusMessage = new Message(Encoding.UTF8.GetBytes(messageBody)) { // 这里用交易号作为MessageId,你也可以换成其他业务唯一标识 MessageId = tran.TransactionNumber }; await collector.AddAsync(serviceBusMessage); logger.LogInformation($"Received transaction added to servicebus queue with MessageId: {serviceBusMessage.MessageId}"); return new NoContentResult(); } catch (Exception ex) { logger.LogError(StandardEvents.FailureEvent, ex, "Unhandled exception"); } throw new ArgumentException("Bad Request"); }
情况2:使用新版Service Bus扩展(v6+,基于Azure.Messaging.ServiceBus SDK)
如果你的项目用的是v6及以上版本的Service Bus扩展包,需要用ServiceBusMessage类型:
- 把函数参数里的
IAsyncCollector<string>替换为IAsyncCollector<ServiceBusMessage> - 构造
ServiceBusMessage实例时设置MessageId属性
修改后的代码示例:
// 记得引用Azure.Messaging.ServiceBus命名空间 using Azure.Messaging.ServiceBus; [FunctionName(nameof(ReceiveTransaction))] [OpenApiOperation(nameof(ReceiveTransaction), "Transaction", Description = "Receives a transaction and adds to the Transaction Platform")] [OpenApiRequestBody("application/json", typeof(string), Description = "Parameters", Required = true)] [OpenApiResponseWithoutBody(HttpStatusCode.BadRequest, Description = "Request was invalid and could not be processed")] [OpenApiResponseWithoutBody(HttpStatusCode.OK, Description = "Request was valid - transaction posted to platform")] public async Task<IActionResult> ReceiveTransaction( [HttpTrigger(AuthorizationLevel.Anonymous, "post", Route = "v1/receiveTransaction")] HttpRequest req, // 替换为IAsyncCollector<ServiceBusMessage> [ServiceBus("%Transactions:SB:Transactions:IngressQueue%", Connection = "Transactions:SB:Transactions")] IAsyncCollector<ServiceBusMessage> collector, ILogger logger) { try { using var sr = new StreamReader(req.Body); var rawJson = await sr.ReadToEndAsync(); var tran = JsonConvert.DeserializeObject<RetailTransaction>(rawJson, JsonSerialisationUtils.SerialiserSettings); tran.UpdateExtendedProperties(); tran.SetPlatformReceived(); logger.LogInformation(StandardEvents.SuccessEvent, "Incoming transaction: {transactionNumber} validated", tran.TransactionNumber); // 构造ServiceBusMessage对象并设置MessageId var messageBody = JsonConvert.SerializeObject(tran, JsonSerialisationUtils.SerialiserSettings); var serviceBusMessage = new ServiceBusMessage(messageBody) { MessageId = tran.TransactionNumber }; await collector.AddAsync(serviceBusMessage); logger.LogInformation($"Received transaction added to servicebus queue with MessageId: {serviceBusMessage.MessageId}"); return new NoContentResult(); } catch (Exception ex) { logger.LogError(StandardEvents.FailureEvent, ex, "Unhandled exception"); } throw new ArgumentException("Bad Request"); }
注意事项
- 确保项目安装了对应版本的NuGet包:老版本用
Microsoft.Azure.WebJobs.Extensions.ServiceBus,新版本用Microsoft.Azure.WebJobs.Extensions.ServiceBus(进程内模型)或Microsoft.Azure.Functions.Worker.Extensions.ServiceBus(隔离进程模型) - MessageId建议设置为业务唯一值(比如示例中的交易号),这样Service Bus可以基于该属性实现重复消息检测等功能
内容的提问来源于stack exchange,提问作者Bhav
相关产品推荐
相关产品推荐

