Kafka消费者本地批量队列内存泄漏问题排查求助
非托管内存持续增长问题(基于confluent-kafka-dotnet 1.9.3)
问题现象
- 10个消费者处理单个Kafka Topic,向Topic添加消息批次后,非托管内存立即增长,即使长时间闲置仍持续上升(对应图1)
- 向Topic发送50万条消息并启动服务后,内存增长趋势明显(对应图2)
参数调整尝试
已定位问题源于消费者本地队列,调整了以下参数:
- QueuedMinMessages:librdkafka维护的每个Topic+分区本地队列最小消息数,从默认值100000改为100
- QueuedMaxMessagesKbytes:本地队列预取消息的最大千字节数,从默认值65536改为30000
修改参数并重启服务(Topic仍有50万条待处理消息)后,内存增长速度变慢,但泄漏问题未彻底解决,本地Kafka队列未清理已处理的消息。
消费者代码
private async Task StartConsumer(CancellationToken stoppingToken) { try { using (var consumer = new ConsumerBuilder<string, string>(_consumerConfig) .SetErrorHandler((_, e) => _logger.LogError($"Error: {e.Reason}")) .Build()) { consumer.Subscribe(_topicName); while (!stoppingToken.IsCancellationRequested) { ConsumeResult<string, string> result = null; try { result = consumer.Consume(); if (result == null) continue; var message = result.Message.Value; Console.WriteLine($"Consumed message '{message}' at '{result.TopicPartitionOffset}'"); if (message != null) { T deserializedMessage = JsonConvert.DeserializeObject<T>(message); if (deserializedMessage != null) { var handler = await _managerFactory.CreateHandler(_topicName); await handler.HandleAsync(deserializedMessage, _topicName); } } else { _logger.LogInformation("Processed empty message from Kafka"); } _logger.LogInformation($"Processed message from Kafka"); consumer.Commit(result); } catch (OracleException ex) { _logger.LogError(ex, "OracleException" + '\n' + ex.Message + '\n' + ex.InnerException); ProcessFailureMessage(result.Message); } catch (ConsumeException ex) { _logger.LogError(ex, "ConsumerException" + '\n' + ex.Message + '\n' + ex.InnerException); } catch (Exception ex) { _logger.LogError(ex, "Exception" + '\n' + ex.Message + '\n' + ex.InnerException); } } } } catch (Exception ex) { _logger.LogError(ex, "Kafka connection error"); } }
消费者配置
"RequestTimeoutMs": 60000, "TransactionTimeoutMs": 300000, "SessionTimeoutMs": 300000, "EnableAutoCommit": false, "QueuedMinMessages": 100, "QueuedMaxMessagesKbytes": 30000, "AutoOffsetReset": "Earliest", "AllowAutoCreateTopics": true, "PartitionAssignmentStrategy": "RoundRobin"
补充说明
使用的confluent-kafka-dotnet版本为1.9.3,StartConsumer()以长时任务方式启动:
protected override Task ExecuteAsync(CancellationToken stoppingToken) { for (int i = 0; i < _consumersCount; i++) { Task.Factory.StartNew(() => StartConsumer(stoppingToken), stoppingToken, TaskCreationOptions.LongRunning, TaskScheduler.Default); } return Task.CompletedTask; }
内容的提问来源于stack exchange,提问作者Nordennavic
相关产品推荐
相关产品推荐

