Azure Service Bus消费者随时间推移处理速度变慢问题排查
Azure Service Bus消息处理量随时间下降的问题排查与修复
问题背景
我有多个带订阅的Topic,每个订阅对应一个消费者。随着时间推移,消息处理量逐渐下降。已经配置了maxConcurrentCalls,但只有初始阶段生效。当前使用azure/service-bus npm包,基于官方消息流接收示例做了并发调用和自动完成的修改。一开始所有消费者共用一个连接,后来改成每个消费者独立连接,但问题依旧存在。
原代码
let receiver: any let sbClient: ServiceBusClient try { sbClient = new ServiceBusClient(connectionString); logger.info("Azure connection is establish successfully for topic:", topic) } catch (Err) { logger.error("Error in subscribing events from azure service bus", Err) return } try { receiver = sbClient?.createReceiver(topic, subscriptionName); logger.debug(`Receiver for ${topic} Connected Successfully.`) } catch (err) { sbClient?.close() logger.error(`Error in creating receiver for ${topic}`, err) return null } try { const subscription = receiver.subscribe({ processMessage: async (brokeredMessage: ServiceBusReceivedMessage) => { var input = brokeredMessage.body let result = await processData(input) if (result) { await receiver.completeMessage(brokeredMessage) } else { await receiver.abandonMessage(brokeredMessage) } }, processError: async (args: any) => { logger.error(`Error from source ${args.errorSource} occurred: `, args.error); if (isServiceBusError(args.error)) { switch (args.error.code) { case "MessagingEntityDisabled": case "MessagingEntityNotFound": case "UnauthorizedAccess": logger.error( `An unrecoverable error occurred. Stopping processing. ${args.error.code}`, args.error ); await subscription.close(); break; case "MessageLockLost": logger.error(`Message lock lost for message`, args.error); break; case "ServiceBusy": await delay(1000); break; default: logger.error("Error in processing message", args) } } }, }, { autoCompleteMessages: false, maxConcurrentCalls: 100 }); return receiver } catch (err) { await receiver?.close(); logger.error("Error in subscribing messages", err) return null }
核心问题分析
- 未捕获
processData异常:如果processData抛出未处理的异常,会中断当前消息处理,甚至阻塞后续并发任务,导致有效处理量下降。 - 资源清理不彻底:遇到致命错误时,仅关闭订阅但未关闭接收器和客户端连接,会造成资源泄漏,长期积累影响接收效率。
- 未知错误无恢复逻辑:遇到未定义的Service Bus错误时,没有触发重建接收器的操作,导致消费者挂起不再接收消息。
- 消息锁丢失无细化监控:锁丢失后Service Bus会自动重新入队,但缺乏对应消息ID的日志,无法排查是否是锁时长不足导致的频繁丢锁。
修复后的代码
let receiver: ServiceBusReceiver | null = null; let sbClient: ServiceBusClient | null = null; try { sbClient = new ServiceBusClient(connectionString); logger.info(`Azure连接成功建立,Topic: ${topic}`); } catch (err) { logger.error("Azure Service Bus连接失败", err); return null; } try { receiver = sbClient.createReceiver(topic, subscriptionName); logger.debug(`Topic ${topic}的接收器创建成功`); } catch (err) { await sbClient.close(); logger.error(`创建Topic ${topic}接收器失败`, err); return null; } try { const subscription = receiver.subscribe({ async processMessage(brokeredMessage: ServiceBusReceivedMessage) { try { const input = brokeredMessage.body; const result = await processData(input); if (result) { await receiver.completeMessage(brokeredMessage); } else { await receiver.abandonMessage(brokeredMessage); } } catch (processErr) { logger.error(`处理消息失败,消息ID: ${brokeredMessage.messageId}`, processErr); // 异常时放弃消息,让其重新进入队列等待处理 await receiver.abandonMessage(brokeredMessage); } }, async processError(args: ProcessErrorArgs) { logger.error(`错误来源: ${args.errorSource},错误信息:`, args.error); if (isServiceBusError(args.error)) { switch (args.error.code) { case "MessagingEntityDisabled": case "MessagingEntityNotFound": case "UnauthorizedAccess": logger.error(`不可恢复错误,停止处理,错误码: ${args.error.code}`, args.error); // 彻底清理所有资源 await subscription.close(); await receiver.close(); await sbClient.close(); break; case "MessageLockLost": logger.error(`消息锁丢失,消息ID: ${args.error.messageId || '未知'}`, args.error); // 锁丢失后Service Bus会自动将消息重新入队,无需额外操作 break; case "ServiceBusy": logger.warn("Service Bus繁忙,延迟1秒后重试"); await new Promise(resolve => setTimeout(resolve, 1000)); break; default: logger.error("处理消息时发生未知错误", args.error); // 清理现有资源,可在此添加延迟重建逻辑避免服务中断 await subscription.close(); await receiver.close(); await sbClient.close(); break; } } }, }, { autoCompleteMessages: false, maxConcurrentCalls: 100 }); // 监听订阅关闭事件,确保资源全部清理 subscription.on("close", async () => { await receiver?.close(); await sbClient?.close(); logger.info(`Topic ${topic}的订阅已关闭,资源已清理`); }); return receiver; } catch (err) { await receiver?.close(); await sbClient?.close(); logger.error("订阅消息失败", err); return null; }
额外优化建议
- 监控处理时长:给
processData添加耗时日志,排查是否存在慢处理任务阻塞并发。 - 调整消息锁时长:如果
processData处理时间超过30秒,需在Topic订阅配置中延长锁时长,避免频繁丢锁。 - 添加自动重建逻辑:在未知错误分支中,添加延迟重建接收器的代码,避免消费者完全停止工作。
- 配置连接池:通过
ServiceBusClient的连接池配置(如maxConcurrentSessions)优化资源利用,减少连接创建开销。
内容的提问来源于stack exchange,提问作者Priya
相关产品推荐
相关产品推荐

