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

使用LiteDB的C#服务器出现死锁及事务超时问题求助

问题:LiteDB数据库锁超时与死锁排查

场景与问题

  • 用C#开发了基于LiteDB的消息存储服务器,启动时异步GC任务定期删除过期消息,依赖LiteDB线程安全特性未加锁。
  • 测试场景:WPF中100个客户端对应独立任务,互相发送消息并接收回复,测试时GC任务每秒执行一次(实际场景为每周一次)。
  • 问题表现:前1-3秒运行正常,之后GC任务停止,每秒仅能发送一条消息,1分钟后抛出锁超时异常:

Exception thrown: 'LiteDB.LiteException' in LiteDB.dll
Database lock timeout when entering in transaction mode after 00:01:00

  • 推测系统存在死锁问题,以下是简化后的服务器及数据库操作代码:

服务器代码

public class Worker : BackgroundService
{
    private static readonly TimeSpan _collectionInterval = TimeSpan.FromSeconds(1);
    private static readonly TimeSpan _expirationThreshold = TimeSpan.FromMinutes(2);
    
    private readonly ConcurrentDictionary<string, Client> Clients = new();

    private readonly DBHandler _dbHandler = DBHandler.GetInstance(AppDomain.CurrentDomain.BaseDirectory);

    protected override async Task ExecuteAsync(CancellationToken cancellationToken)
    {
        // 创建关联的CTSource,以便GC任务在致命异常时取消整个操作
        CancellationTokenSource cancellationTokenSource =
            CancellationTokenSource
            .CreateLinkedTokenSource(cancellationToken);
        CancellationToken token = cancellationTokenSource.Token;

        List<Task> GCAndConnectionsTasks = new();

        try
        {
            listener.Start();

            // 启动垃圾回收任务,定期从DB中移除过期消息
            // 在服务器生命周期内每隔_collectionInterval执行一次
            GCAndConnectionsTasks.Add(
                CollectExpiredMessagesAsync(
                    collectionInterval: _collectionInterval,
                    expirationThreshold: _expirationThreshold,
                    cancellationTokenSource // 允许任务取消服务器
                )
            );

            while (!token.IsCancellationRequested)
            {
                TcpClient tcpClient = await listener.AcceptTcpClientAsync(token);

                // 在独立任务中处理客户端连接
                GCAndConnectionsTasks.Add(
                    HandleClientConnectionAsync(tcpClient, token)
                );
            }

            cancellationTokenSource.Cancel();
            await Task.WhenAll(GCAndConnectionsTasks);
        }
        finally
        {
            _dbHandler.Dispose();
            listener.Stop();
            cancellationTokenSource.Dispose();
            token.ThrowIfCancellationRequested();
        }
    }

    private async Task HandleClientConnectionAsync(TcpClient tcpClient, CancellationToken token)
    {
        Client? client = null;

        try
        {
            client = new(tcpClient);
            client.EnableKeepAlive();

            await SendMessageAsync(new WelcomeMessage(), client, token);

            DisconnectReason? reason = null;
            Message? message = await ReceiveMessageAsync(client, token);
            
            if (Clients.ContainsKey(message.FromStationID))
            {
                Clients.TryAdd(message.FromStationID, client);
                await ReceiveMessagesAsync(client, token);
            }
        }
        finally
        {
            // 清理已关闭的连接
            if (client is not null)
            {
                if (client.StationID is not null)
                {
                    Clients.TryRemove(client.StationID, out _);
                }
                client.CloseTcpAndStream();
            }
        }
    }

    private async Task ReceiveMessagesAsync(Client client, CancellationToken token)
    {
        while (client.IsConnected && !token.IsCancellationRequested)
        {
            Message? message = await ReceiveMessageAsync(client, token);
            
            if (token.IsCancellationRequested || message is null)
            {
                break;
            }

            await ProcessMessageByTypeAsync(message, token);
        }
    }

    private async Task ProcessMessageByTypeAsync(Message message, CancellationToken token)
    {
        if (message is null)
        {
            return;
        }
        else if (message is AckMessage ackMessage)
        {
            await ProcessAckMessageAsync(ackMessage, token);
        }
        else if (message is ComponentsMessage || message is FeedbackMessage)
        {
            await ProcessDataMessageAsync(message, token);
        }
        // 忽略其他消息类型
    }

    private async Task ProcessDataMessageAsync(Message message, CancellationToken token)
    {
        if (message is null)
        {
            return;
        }

        if (message.ToStationID != null)
        {
            _dbHandler.Insert(message);

            Client client;

            if (Clients.TryGetValue(message.ToStationID, out client))
            {
                await SendMessageAsync(message, client, token);
            }
        }
    }

    private async Task ProcessAckMessageAsync(AckMessage ackMessage, CancellationToken token)
    {
        _dbHandler.DeleteMessageByID(ackMessage.AckedMessageID);
    }

    private static async Task SendMessageAsync(Message message, Client client, CancellationToken token)
    {
        // 使用message.Serialize(),然后添加4字节头部生成byte[] buffer
        
        await client.WriteAsync(buffer, token);
    }

    private static async Task<Message?> ReceiveMessageAsync(Client client, CancellationToken token)
    {
        // 使用await client.ReadAsync接收4字节头部
        // 使用await client.ReadAsync接收消息字节到byte[] buffer
        
        return
            buffer is null ?
            null :
            Message.Deserialize(Encoding.UTF8.GetString(buffer, 0, buffer.Length));
    }

    private async Task CollectExpiredMessagesAsync(
        TimeSpan collectionInterval,
        TimeSpan expirationThreshold,
        CancellationTokenSource source)
    {
        while (!source.IsCancellationRequested)
        {
            // 确保GC对DB的Messages集合拥有独占访问权
            
            try
            {
                _dbHandler.DeleteExpiredMessages(expirationThreshold);
            }
            // 致命错误,取消整个操作(服务器)
            catch (Exception ex)
            {
                source.Cancel();
                throw;
            }
            

            if (!source.IsCancellationRequested)
            {
                // 下次垃圾回收前延迟等待
                await Task.Delay(collectionInterval, source.Token);
            }
        }
    }
}

数据库操作代码

private ILiteCollection<Message> GetCollection(Message message)
{
    var msgCollection = _dataBase
        .GetCollection<Message>(DB.MESSAGES);

    msgCollection.EnsureIndex((Message m) => m.MessageID);

    return msgCollection;
}

private ILiteCollection<Message> GetMessages()
{
    return GetCollection((Message)null);
}

public BsonValue Insert(Message message)
{
    message.MessageID = 0;
    message.TimeStamp = DateTime.UtcNow;
    return GetMessages()
        .Insert(message);
}

public bool DeleteMessageByID(BsonValue messageID)
{
    return GetMessages()
        .Delete(messageID);
}

public IEnumerable<Message> GetMessagesForStationID(string stationID)
{
    return GetMessages()
        .Find(message => message.ToStationID ==  stationID);
}

public int DeleteExpiredMessages(TimeSpan expirationThreshold)
{
    var cutoffTime = DateTime.UtcNow - expirationThreshold;
    return GetMessages()
        .DeleteMany((Message m) => m.TimeStamp < cutoffTime);
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 17:24:54