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

Azure Service Bus消息重发后莫名进入DLQ的异常排查

问题分析与解决方案

核心原因

你的代码存在几个关键问题,直接导致消息锁超时进入死信队列:

1. 未正确释放MessageReceiver资源

每次调用Test方法时,都会创建新的MessageReceiver实例并覆盖静态变量_messageReceiver,但旧的实例从未被Dispose。当MessageReceiver未被正确关闭时,它持有的消息锁不会主动释放,直到Service Bus的锁超时时间到期。而你的主题Max delivery count设为1,锁超时后消息会直接进入死信队列。

2. 仅完成未超时的消息

代码中只对LockedUntilUtc > DateTime.UtcNow的消息调用CompleteAsync,但在消息接收、处理(重发)的过程中,部分消息的锁可能已经超时。这些超时的消息不会被完成,锁到期后会触发投递计数增加,因为投递计数是1,直接进入DLQ。

3. 静态MessageReceiver的滥用

静态变量在Azure函数的多个执行实例或调用间共享,会导致Receiver实例混乱:旧实例未关闭,新实例又被创建,多个Receiver同时操作同一订阅,引发锁管理的冲突。

修复方案

针对上述问题,修改代码如下:

1. 用using包裹MessageReceiver,确保资源释放

每次创建MessageReceiver时,使用using语句,保证方法执行完后自动关闭Receiver,释放所有未处理的锁(如果有的话)。

2. 移除锁超时判断,统一完成所有接收的消息

不管消息锁是否超时,都尝试调用CompleteAsync——即使锁已经超时,调用CompleteAsync只会返回失败(不会抛出异常影响流程),但能确保所有被接收的消息都被正确处理,避免锁残留。

3. 避免静态MessageReceiver,改用局部变量

静态Receiver会引发多实例冲突,改用局部变量配合using管理生命周期。

修改后的关键代码片段

public async Task Test(ILogger logger)
{
    // 使用using包裹MessageReceiver,确保自动释放
    using var messageReceiver = new MessageReceiver(
        ConnectionStrings.ServiceBusConnectionString, 
        EntityNameHelper.FormatSubscriptionPath("plandisc.domain.models.servicebus.testresendmessage", "testresendmessage-subscription"), 
        ReceiveMode.PeekLock, 
        null, 
        0);

    logger.Info("Receiving Messages");
    var messages = new List<Message>();
    IList<Message> messageBatch;

    do
    {
        messageBatch = await messageReceiver.ReceiveAsync(int.MaxValue, TimeSpan.FromSeconds(2));
        if (messageBatch != null && messageBatch.Any()) messages.AddRange(messageBatch);
        logger.Info("ReceiveAsync from " + messageReceiver.Path + " returned: " + (messageBatch?.Count() ?? 0));
    } while (messageBatch != null && messageBatch.Any());

    // 消息重发逻辑保持不变
    var serviceBusMessages = messages.Select(m => Json.Deserialize<TestResendMessage>(System.Text.Encoding.UTF8.GetString(m.Body)));
    foreach (var message in serviceBusMessages)
    {
        logger.Info("Resending message: " + message.Data + " " + message.TimeStamp);
        await SendAsync(new TestResendMessage
        {
            Data = message.Data,
            TimeStamp = message.TimeStamp
        }, logger);
    }

    logger.Info($"Complete messages (Messages: {messages.Count})");

    if (messages.Any())
    {
        var info = "messages: " + messages.Count() + "\n" +
            "lockedUntilUtc: " + string.Join(", ", messages.Select(m => m.SystemProperties.LockedUntilUtc.ToString()));
        logger.Info(info);

        // 移除锁超时判断,统一尝试完成所有消息
        var lockTokens = messages.Select(m => m.SystemProperties.LockToken);
        try
        {
            await messageReceiver.CompleteAsync(lockTokens);
        }
        catch (ServiceBusException ex)
        {
            // 捕获锁超时等异常,记录日志但不中断流程
            logger.Error($"Failed to complete some messages: {ex.Message}", ex);
        }
    }

    logger.Info("Done completing messages");
}

额外优化建议

  • 调整主题的Max delivery count为大于1的值(比如3),给消息重试的机会,避免一次锁超时就进入DLQ。
  • 监控Service Bus的锁超时时间(默认是60秒),如果你的消息处理(重发)耗时较长,在处理过程中调用messageReceiver.RenewLockAsync延长锁时间。
  • 静态TopicClient可以保留,但建议改用ServiceBusClient(新的Azure.Messaging.ServiceBus SDK)替代旧的Microsoft.Azure.ServiceBus SDK——旧SDK已停止维护,新SDK有更好的连接管理和稳定性。

内容的提问来源于stack exchange,提问作者jbiversen

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 09:36:23