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

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.18 15:50:39