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

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
}

核心问题分析

  1. 未捕获processData异常:如果processData抛出未处理的异常,会中断当前消息处理,甚至阻塞后续并发任务,导致有效处理量下降。
  2. 资源清理不彻底:遇到致命错误时,仅关闭订阅但未关闭接收器和客户端连接,会造成资源泄漏,长期积累影响接收效率。
  3. 未知错误无恢复逻辑:遇到未定义的Service Bus错误时,没有触发重建接收器的操作,导致消费者挂起不再接收消息。
  4. 消息锁丢失无细化监控:锁丢失后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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 23:54:59