如何在Azure函数中重调度Service Bus队列消息并递增deliveryCount
我使用Azure Service Bus Queue收集需要发送至API的消息,并借助Queue上的Azure Trigger Function处理入站消息。但该API不可靠,因此需要添加自定义重试计划:首次处理失败时15分钟后重试,之后依次为1小时、2小时……当达到maxDeliveryCount时,消息应被发送至DLQ。
目前用@azure/service-bus包实现的重调度功能正常,但每次调度的消息都会被视为新消息,导致deliveryCount始终为1。我不想在消息体中添加自定义deliveryCount属性污染元数据,因此想问:如何在重调度时递增原生的deliveryCount?
最小可复现示例(MRE)
const serviceBusQueueTrigger: AzureFunction = async function ( context: Context, msg: any ): Promise<void> { context.log("ServiceBus queue trigger function processed message", msg); const serviceBusClient = new ServiceBusClient(<connection_string>); const sender = serviceBusClient.createSender(<queue_name>); context.log("deliveryCount: ", context.bindingData.deliveryCount); try { // POST to API } catch (error) { // 若此处抛出错误,触发器会立即重试,我希望根据当前deliveryCount设置延迟 await sender.scheduleMessages( { body: msgData, contentType: "application/json", }, // 示例延迟时间 new Date(Date.now() + Math.pow(2, context.bindingData.deliveryCount) * 60 * 1000) ); } await sender.close(); }
编辑说明:尽管微软文档仍有内置重试机制的模糊说明,但该功能已不再受支持。初始问题仍然有效:如何用Azure Service Bus Trigger Function构建该自定义重试机制?
当前配置文件
host.json
{ "version": "2.0", "logging": { "applicationInsights": { "samplingSettings": { "isEnabled": true, "excludedTypes": "Request" } } }, "extensionBundle": { "id": "Microsoft.Azure.Functions.ExtensionBundle", "version": "[4.0.0, 5.0.0)" }, "concurrency": { "dynamicConcurrencyEnabled": true, "snapshotPersistenceEnabled": true }, "extensions": { "serviceBus": { "clientRetryOptions": { "mode": "exponential", "tryTimeout": "00:01:00", "delay": "00:01:00", "maxDelay": "03:00:00", "maxRetries": 3 } } } }
function.json
{ "bindings": [ { "name": "incoming", "type": "serviceBusTrigger", "direction": "in", "queueName": "incoming", "connection": "SERVICEBUS" } ], "retry": { "strategy": "exponentialBackoff", "maxRetryCount": 5, "minimumInterval": "00:00:10", "maximumInterval": "00:15:00" }, "scriptFile": "../dist/trigger/index.js" }
不要手动创建新消息进行调度,而是利用Service Bus的延迟原消息特性,保留消息原生元数据并自动递增deliveryCount:
1. 修改函数逻辑,延迟原消息而非创建新消息
直接对当前接收的消息执行延迟操作,这样消息会保留原有的deliveryCount,延迟到期后重新投递时deliveryCount自动+1:
const serviceBusQueueTrigger: AzureFunction = async function ( context: Context, msg: any, messageReceiver: ServiceBusReceiver // 通过绑定获取消息接收器 ): Promise<void> { const currentDeliveryCount = context.bindingData.deliveryCount; context.log("当前deliveryCount: ", currentDeliveryCount); try { // 调用API逻辑 await context.done(); // 处理成功,完成消息确认 } catch (error) { // 按自定义规则计算延迟时间 let delayMs: number; switch(currentDeliveryCount) { case 1: delayMs = 15 * 60 * 1000; // 首次失败延迟15分钟 break; case 2: delayMs = 60 * 60 * 1000; // 第二次延迟1小时 break; case 3: delayMs = 2 * 60 * 60 * 1000; // 第三次及之后延迟2小时 break; default: delayMs = 2 * 60 * 60 * 1000; } // 延迟当前消息,保留原生元数据 await messageReceiver.scheduleMessageDelivery(msg, new Date(Date.now() + delayMs)); } }
2. 配置Queue的maxDeliveryCount
在Azure Portal的Service Bus Queue设置中,将maxDeliveryCount设为你期望的最大重试次数(例如设为4,对应首次投递+3次重试)。当消息的deliveryCount达到该值时,Service Bus会自动将其移入DLQ,无需手动处理。
3. 调整函数绑定与配置
- 确保Service Bus Trigger使用默认的
peekLock接收模式(默认配置即为该模式),这样才能对消息执行延迟操作。 - 移除或调整host.json中的
clientRetryOptions和function.json中的retry配置,避免与自定义重试逻辑冲突。
关键原理
- 手动调度新消息会生成全新的消息实例,
deliveryCount重置为1;而延迟原消息会保留其所有原生元数据,每次重新投递时deliveryCount自动递增。 - 当
deliveryCount达到Queue的maxDeliveryCount阈值时,Service Bus会自动将消息移入DLQ,无需额外代码处理。
内容的提问来源于stack exchange,提问作者picklepick

