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

LiteDB死锁问题排查:C#高并发服务器数据库超时异常

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

问题描述

我用C#编写服务器,基于LiteDB存储消息。服务器收到消息后存入数据库并转发至指定目标。因为LiteDB宣称线程安全,所以未加锁同步读写操作。

测试场景:WPF应用启动100个独立任务作为客户端,每个客户端向其他所有客户端发送消息,同时给客户端和消息添加随机延迟。测试异常表现:前400-2000条消息能正常处理,之后处理速度骤降到每秒1条,最终等待1分钟后抛出异常:

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

无消息延迟时问题出现更早,设置约50ms延迟则可通过测试。怀疑是LiteDB内部死锁或高并发任务导致。

服务器核心代码

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)
    {
        CancellationTokenSource cancellationTokenSource = CancellationTokenSource
                .CreateLinkedTokenSource(cancellationToken);
        CancellationToken token = cancellationTokenSource.Token;
    
        List<Task> tasks = new();

        try
        {
            listener.Start();

            while (!token.IsCancellationRequested)
            {
                TcpClient tcpClient;
                try
                {
                    // Accept an incoming client connection
                    tcpClient = await listener.AcceptTcpClientAsync(token);
                }
                catch (Exception ex)
                {
                    cancellationTokenSource.Cancel();
                    break;
                }

                // Handle the client connection in a separate task
                tasks.Add(
                    HandleClientConnectionAsync(tcpClient, token)
                );
            }

            cancellationTokenSource.Cancel();
            await Task.WhenAll(tasks);
        }
        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);

            Message? message = await ReceiveMessageAsync(client, token);
            
            if (Clients.TryAdd(message.FromStationID, client))
            {
                await ReceiveMessagesAsync(client, token);
            }
        }
        finally
        {
            // Clean up closed connection
            if (client is not null)
            {
                if (client.StationID is not null)
                {
                    Clients.TryRemove(client.StationID, out _);
                }
                client.CloseTcpAndStream();
            }
        }
    }
    
    private static async Task SendMessageAsync(Message message, Client client, CancellationToken token)
    {
        //use message.Serialize(), then add 4 byte header to create byte[] buffer
        
        await client.WriteAsync(buffer, token);
    }
    
    private static async Task<Message?> ReceiveMessageAsync(Client client, CancellationToken token)
    {
        // use await client.ReadAsync to receive 4 byte header
        // use await client.ReadAsync to receive the message bytes into byte[] buffer
        
        return
            buffer is null ?
            null :
            Message.Deserialize(Encoding.UTF8.GetString(buffer, 0, buffer.Length));
    }

    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 DataMessage || message is FeedbackMessage)
        {
            await ProcessDataMessageAsync(message, token);
        }
        // Ignore other messages
    }

    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 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;
    return GetMessages()
        .Insert(message);
}

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

问题根源排查与解决方案

1. LiteDB的"线程安全"并非无限制高并发安全

LiteDB的线程安全是指多线程/任务可以安全调用API,但它采用单写多读的锁机制:写操作会排他性锁定整个数据库,读操作共享锁。100个客户端高并发写入时,所有写请求会排队等待锁,队列积压后会出现处理速度骤降,最终触发超时。

2. EnsureIndex重复调用引发锁竞争

在GetCollection方法中,每次插入/删除都会调用EnsureIndex。虽然该方法内部会检查索引是否存在,但检查+创建的过程需要加锁,高并发下频繁调用会导致大量锁竞争,加剧阻塞。

修复: 将EnsureIndex移到数据库初始化阶段,仅调用一次:

// 在DBHandler的构造函数或初始化方法中执行
_dataBase.GetCollection<Message>(DB.MESSAGES).EnsureIndex(m => m.MessageID);

3. 同步数据库操作阻塞异步线程

Insert和DeleteMessageByID是同步方法,在异步消息处理流程中调用会阻塞当前线程。高并发下会耗尽线程池资源,导致后续任务无法调度,进一步加剧锁等待。

修复: 改用LiteDB的异步API:

// 修改Insert为异步
public async Task<BsonValue> InsertAsync(Message message)
{
    message.MessageID = 0;
    return await _dataBase.GetCollection<Message>(DB.MESSAGES).InsertAsync(message);
}

// 修改Delete为异步
public async Task<bool> DeleteMessageByIDAsync(BsonValue messageID)
{
    return await _dataBase.GetCollection<Message>(DB.MESSAGES).DeleteAsync(messageID);
}

在消息处理方法中调用异步版本:

private async Task ProcessDataMessageAsync(Message message, CancellationToken token)
{
    if (message is null || message.ToStationID == null) return;
    
    await _dbHandler.InsertAsync(message);

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

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

4. 事务超时设置与锁等待队列

LiteDB默认事务超时为1分钟,高并发下锁等待队列过长会触发超时。可按需调整超时时间(仅缓解,核心仍需减少锁竞争):

// 初始化数据库时设置超时(示例:2分钟)
var db = new LiteDatabase(@"Filename=MyData.db;Timeout=120");

5. 数据库实例管理

确保DBHandler是单例且正确释放,检查_dataBase是否在Dispose方法中正确调用Dispose,避免连接泄露导致锁问题。

6. 高并发写入限流

如果业务允许,可对数据库写入操作做限流,用SemaphoreSlim控制并发写入数量,避免瞬间大量请求压垮LiteDB:

// 在Worker类中定义限流信号量
private readonly SemaphoreSlim _dbWriteSemaphore = new SemaphoreSlim(10); // 允许10个并发写入

// 在ProcessDataMessageAsync中使用
await _dbWriteSemaphore.WaitAsync(token);
try
{
    await _dbHandler.InsertAsync(message);
}
finally
{
    _dbWriteSemaphore.Release();
}

内容的提问来源于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:56:00