使用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
相关产品推荐
相关产品推荐

