LiteDB死锁问题排查:C#高并发服务器数据库超时异常
问题描述
我用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

