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

使用ServiceBusSessionProcessor跨线程完成消息报错问题咨询

问题分析与解决方案

问题根源

你遇到的ServiceBusReceiver has already been closed错误,核心原因有两个:

  1. SessionIdleTimeout与批量处理窗口不匹配:你设置了SessionIdleTimeout = TimeSpan.FromSeconds(2),但要等5秒才触发批量处理。当session在2秒内没有新消息时,SessionProcessor会自动关闭对应的Receiver并释放session,此时再用之前保存的ProcessSessionMessageEventArgs调用CompleteMessageAsync自然会报错。
  2. 错误保存ProcessSessionMessageEventArgs:该对象与当前session的Receiver强绑定,session释放后,关联的Receiver资源会被回收,无法再用于操作消息。

修正方案

1. 对齐SessionIdleTimeout与批量处理时间

将SessionIdleTimeout设置为大于等于批量处理的等待时间(比如你设置的5秒),确保在批量处理完成前,session不会被自动释放:

processor = _client.CreateSessionProcessor("mikes.test.queue.session", new ServiceBusSessionProcessorOptions()
{
    ReceiveMode = ServiceBusReceiveMode.PeekLock,
    PrefetchCount = prefetchCount,
    AutoCompleteMessages = autocompleteMessage,
    SessionIdleTimeout = TimeSpan.FromSeconds(10), // 设为比5秒长的时间
    MaxConcurrentSessions = concurrentSessionsPerProcessor
});

2. 保存消息关键信息而非EventArgs

不要直接存储ProcessSessionMessageEventArgs,而是提取消息的LockToken、SessionId等必要信息,后续通过ServiceBusClient创建新的Receiver来完成消息:

// 替换原GroupedSession类
public class GroupedSession
{
    public DateTime Timestamp { get; set; }
    public List<(string SessionId, string LockToken)> Messages { get; set; } = new List<(string, string)>();
}

// 修改MessageHandler
async Task MessageHandler(ProcessSessionMessageEventArgs args)
{
    string sessionId = args.Message.SessionId;
    string lockToken = args.Message.LockToken;

    lock (_lock)
    {
        if (!_messages.ContainsKey(sessionId))
            _messages.Add(sessionId, new GroupedSession() 
            { 
                Timestamp = DateTime.Now,
                Messages = new List<(string, string)>() { (sessionId, lockToken) }
            });
        else
            _messages[sessionId].Messages.Add((sessionId, lockToken));
    }

    await Task.CompletedTask;
}

3. 优化批量处理逻辑

  • 使用异步方式处理批量操作,避免阻塞线程;
  • 遍历字典时先复制键的列表,防止遍历过程中字典修改引发异常;
  • 用ServiceBusReceiver批量完成消息,提升效率:
// 替换原ProcessMessageTimerTick和CompleteMessageList方法
private async void ProcessMessageTimerTick(object source, ElapsedEventArgs e)
{
    // 先复制所有sessionId,避免遍历中字典变化
    List<string> sessionIds;
    lock (_lock)
    {
        sessionIds = _messages.Keys.ToList();
    }

    foreach (var sessionId in sessionIds)
    {
        GroupedSession groupedSession;
        lock (_lock)
        {
            if (!_messages.TryGetValue(sessionId, out groupedSession))
                continue;
            if (groupedSession.Timestamp.AddSeconds(5) >= DateTime.Now)
                continue;
            
            // 取出并移除该session的消息
            var messages = groupedSession.Messages.ToList();
            _messages.Remove(sessionId);
        }

        // 批量完成消息
        await CompleteMessagesInBatch(sessionId, messages.Select(m => m.LockToken).ToList());
    }
}

private async Task CompleteMessagesInBatch(string sessionId, List<string> lockTokens)
{
    if (lockTokens.Count == 0)
        return;

    // 创建针对该session的Receiver
    using var receiver = _client.CreateReceiver("mikes.test.queue.session", new ServiceBusReceiverOptions
    {
        SessionId = sessionId,
        ReceiveMode = ServiceBusReceiveMode.PeekLock
    });

    try
    {
        await receiver.CompleteMessagesAsync(lockTokens);
    }
    catch (Exception ex)
    {
        // 处理错误,比如记录日志、重试或死信
        Console.WriteLine($"完成session {sessionId}消息失败: {ex.Message}");
    }
}

4. 避免同步阻塞异步方法

原代码中CompleteMessageList(array).GetAwaiter().GetResult()会阻塞线程,容易引发死锁,改用异步方式处理(如上面的async void定时器回调,注意仅在定时器场景使用)。

额外注意事项

  • 消息锁过期问题:如果批量处理时间超过消息的LockDuration(默认30秒),需要提前调用RenewMessageLockAsync续约锁。可以在定时器中定期对未处理的消息续约;
  • 线程安全:确保字典的访问始终加锁,避免并发问题;
  • 资源释放:创建的ServiceBusReceiver要使用using语句自动释放资源。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 22:13:14