QuestDB ILP连接被远程主机强制关闭问题求助
问题描述
通过ILP向QuestDB写入数据时,应用偶尔停止写入并抛出以下错误:
System.IO.IOException: Unable to write data to the transport connection: An existing connection was forcibly closed by the remote host.. ---> System.Net.Sockets.SocketException (10054): An existing connection was forcibly closed by the remote host. at System.Net.Sockets.Socket.AwaitableSocketAsyncEventArgs.CreateException(SocketError error, Boolean forAsyncThrow) at System.Net.Sockets.Socket.AwaitableSocketAsyncEventArgs.SendAsyncForNetworkStream(Socket socket, CancellationToken cancellationToken) at System.Net.Sockets.NetworkStream.WriteAsync(Byte[] buffer, Int32 offset, Int32 count, CancellationToken cancellationToken) at QuestDB.LineTcpSender.SendAsync(CancellationToken cancellationToken) at System.Runtime.CompilerServices.AsyncMethodBuilderCore.Start[TStateMachine](TStateMachine& stateMachine) at QuestDB.LineTcpSender.SendAsync(CancellationToken cancellationToken) at nL1nO95G5COwGUKOWaD1.SC9pfx5voyF(Object, CancellationToken, nL1nO95G5COwGUKOWaD1) at BMPolytec.AppServerKMR.Services.QuestDbService.Blgaq8dfyv() at System.Runtime.CompilerServices.AsyncTaskMethodBuilder1.AsyncStateMachineBox1.ExecutionContextCallback(Object s) at System.Threading.ExecutionContext.RunInternal(ExecutionContext executionContext, ContextCallback callback, Object state)
环境信息:
- 3个ILP连接分别对应写入3张表,每分钟总计写入60行数据
- 4台配置完全相同的PC运行同款软件,仅3台出现该错误;第4台的区别是每100ms插入60行,而非每分钟60行
- QuestDB版本:7.1.3
相关代码片段:
SpeedSender = await LineTcpSender.ConnectAsync(Configuration.IpAddress, Configuration.Port, tlsMode: TlsMode.Disable); while (true) { try { if (SpeedQueue.Count > 0) { while (SpeedQueue.Count != 0) { SpeedClass speed; lock (_speedLock) speed = SpeedQueue.Dequeue(); SpeedSender.Table("Speed") .Column("Speed", speed.Speed) .Column("Type", (int)speed.SpeedType) .At(DateTime.UtcNow); await SpeedSender.SendAsync(); } } } catch (Exception e) { _log.Error(e, $"Dequeue speed"); } Thread.Sleep(50); }
原因分析
- 连接空闲超时:QuestDB的ILP服务默认会断开超过30秒无数据传输的空闲连接。3台低频率写入的机器,队列处理完毕后可能出现超过30秒的空闲窗口,触发QuestDB的连接回收机制;而第4台高频率写入的机器连接始终有数据交互,不会触发超时。
- 单条发送的网络开销:每处理一行就调用一次
SendAsync,会产生大量TCP交互,不仅效率低下,还增加了连接因网络波动或服务端策略被断开的概率。
解决方案
1. 增加连接重连逻辑
当前代码捕获异常仅记录日志,未尝试重建连接,导致后续写入完全中断。需在异常处理中重新初始化发送器:
catch (Exception e) { _log.Error(e, $"Dequeue speed, reconnecting..."); // 清理旧连接 try { SpeedSender?.Dispose(); } catch (Exception disposeEx) { _log.Warn(disposeEx, $"Failed to dispose old sender"); } // 重建连接 SpeedSender = await LineTcpSender.ConnectAsync(Configuration.IpAddress, Configuration.Port, tlsMode: TlsMode.Disable); }
2. 批量发送数据
将多条数据累积后一次性发送,减少网络交互次数:
SpeedSender = await LineTcpSender.ConnectAsync(Configuration.IpAddress, Configuration.Port, tlsMode: TlsMode.Disable); while (true) { try { if (SpeedQueue.Count > 0) { // 每次批量处理最多50条,可根据实际调整 int batchSize = Math.Min(SpeedQueue.Count, 50); for (int i = 0; i < batchSize; i++) { SpeedClass speed; lock (_speedLock) speed = SpeedQueue.Dequeue(); SpeedSender.Table("Speed") .Column("Speed", speed.Speed) .Column("Type", (int)speed.SpeedType) .At(DateTime.UtcNow); } // 批量发送 await SpeedSender.SendAsync(); } } catch (Exception e) { _log.Error(e, $"Dequeue speed, reconnecting..."); try { SpeedSender?.Dispose(); } catch (Exception disposeEx) { _log.Warn(disposeEx, $"Failed to dispose old sender"); } SpeedSender = await LineTcpSender.ConnectAsync(Configuration.IpAddress, Configuration.Port, tlsMode: TlsMode.Disable); } Thread.Sleep(50); }
3. 调整QuestDB空闲超时(可选)
如果必须保持低频率写入,可修改QuestDB配置文件server.conf,延长ILP连接的空闲超时:
line.tcp.idle.timeout.sec=300 # 改为300秒,根据需求调整
修改后需重启QuestDB生效。
4. 优化队列处理逻辑
使用BlockingCollection替代普通队列,简化生产者-消费者模型,避免手动锁和轮询:
private BlockingCollection<SpeedClass> _speedQueue = new BlockingCollection<SpeedClass>(); // 初始化发送器 SpeedSender = await LineTcpSender.ConnectAsync(Configuration.IpAddress, Configuration.Port, tlsMode: TlsMode.Disable); while (!_speedQueue.IsCompleted) { try { List<SpeedClass> batch = new List<SpeedClass>(); // 批量获取数据,超时50ms while (_speedQueue.TryTake(out var speed, 50)) { batch.Add(speed); if (batch.Count >= 50) break; } if (batch.Count > 0) { foreach (var speed in batch) { SpeedSender.Table("Speed") .Column("Speed", speed.Speed) .Column("Type", (int)speed.SpeedType) .At(DateTime.UtcNow); } await SpeedSender.SendAsync(); } } catch (Exception e) { _log.Error(e, $"Processing speed queue, reconnecting..."); try { SpeedSender?.Dispose(); } catch (Exception disposeEx) { _log.Warn(disposeEx, $"Failed to dispose old sender"); } SpeedSender = await LineTcpSender.ConnectAsync(Configuration.IpAddress, Configuration.Port, tlsMode: TlsMode.Disable); } }
内容的提问来源于stack exchange,提问作者michele bridi
相关产品推荐
相关产品推荐

