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

Java应用使用Kafka Producer出现大量TCP-ESTABLISHED连接的原因排查

Kafka Producer TCP连接数持续增长问题排查与解决

我在Java应用中使用apache kafka-client-3.2.0(也曾尝试最新版apache kafka-client-3.4.0)创建Producer。仅部署1个Broker,创建了5个topic,每个topic的replication-factor=1、partitions=1。其中1个topic由log4j2的KafkaAppender使用,另外4个由定时任务周期性调用Producer写入数据。启动应用后,TCP连接数持续增长,已超过300个且仍在增加。使用过的Kafka版本包括kafka-2.13_2.8.0、kafka-2.13_3.3.1和kafka-2.13_3.4.0。

现有Producer配置代码

private KafkaProducer<String, String> producer;
public KafkaProducer<String, String> getProducer() {
    return producer;
}
private KafkaProducer<String, String> createAndGetProducer(String acksConfig, int retryConfig){
    Properties props = new Properties();
    props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, brokers);
    props.put(ProducerConfig.ACKS_CONFIG, acksConfig);
    props.put(ProducerConfig.RETRIES_CONFIG, retryConfig);
    props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, 
        "org.apache.kafka.common.serialization.StringSerializer");
    props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, 
        "org.apache.kafka.common.serialization.StringSerializer");
    KafkaProducer<String, String> producer = new KafkaProducer<>(props);
    return producer;
}

定时任务写Topic代码

public static void writeToTopic(String topicName, String value){
    logger.info("Writing into topic " + topicName);
    
    ProducerRecord <String, String> producerData = new ProducerRecord <String, String> (topicName, value);
    KafkaProducer<String, String> producer = KafkaConnectionManager.getConnection().getProducer();
    
    logger.info(producer.toString());
    KafkaConnectionManager.getConnection().getProducer().send(producerData);
    KafkaConnectionManager.getConnection().getProducer().flush();
}

当前TCP连接信息

tcp6       0      0  X.X.X.X:9092       X.X.X.X:59604      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:48358      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:50668      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:59898      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:50626      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:59444      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:49716      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:61049      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:58516      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:51514      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:51892      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:50802      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:47106      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:60084      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:50588      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:50682      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:59814      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:50788      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:48554      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:51390      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:58576      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:50974      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:51504      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:47262      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:60022      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:50558      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:56118      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:51844      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:59712      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:51834      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:50626      ESTABLISHED 4369/java
tcp6       0      0  X.X.X.X:9092       X.X.X.X:59668      ESTABLISHED 4369/java

Kafka连接管理器代码

private static KafkaConnectionManager kafkaConnectionManager = null;

public synchronized static KafkaConnectionManager getConnection() {     
    return kafkaConnectionManager;
}

public synchronized static KafkaConnectionManager initialize(Properties kafkaProps, Logger logger, String hostIp) throws KafkaConnectionManagerException {      
    if (kafkaConnectionManager == null) {
        synchronized (KafkaConnectionManager.class) {
            if (kafkaConnectionManager == null) {
                kafkaConnectionManager = new KafkaConnectionManager(kafkaProps, logger, hostIp);
            }
        }
    }
    return kafkaConnectionManager;
}

private KafkaConnectionManager(Properties kafkaProps, Logger logger, String hostIp) throws KafkaConnectionManagerException{
    if(kafkaProps == null)
    {
        throw new KafkaConnectionManagerException("kafkaProps property is null");
    }
    else if(logger == null)
    {
        throw new KafkaConnectionManagerException("logger is null");
    }
    else if(hostIp == null)
    {
        throw new KafkaConnectionManagerException("hostIp is null, set the hostIp");
    }
    KafkaConnectionManager.kafkaProps = kafkaProps;
    this.logger = logger;
    this.hostIp = hostIp;
    brokers = kafkaProps.getProperty(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG).trim();
    logger.info("KAFKA brokers!!! = "+brokers);
    this.acksConfig = kafkaProps.getProperty(ProducerConfig.ACKS_CONFIG).trim();
    this.retryConfig = Integer.parseInt(kafkaProps.getProperty(ProducerConfig.RETRIES_CONFIG).trim());
    
    this.producer = createAndGetProducer(acksConfig, retryConfig);
    logger.info("Kafka Producer has been Initialized successfully");
}

注:initialize方法在main方法中调用

Broker配置

broker.id=1
listeners=PLAINTEXT://:9092
advertised.listeners=PLAINTEXT://X.X.X.X:9092
listener.security.protocol.map=PLAINTEXT:PLAINTEXT,SSL:SSL,SASL_PLAINTEXT:SASL_PLAINTEXT,SASL_SSL:SASL_SSL
num.network.threads=3 
num.io.threads=8 
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400 
socket.request.max.bytes=104857600

log.dirs=/tmp/kafka-logs
num.partitions=1
num.recovery.threads.per.data.dir=1

offsets.topic.replication.factor=1 transaction.state.log.replication.factor=1
transaction.state.log.min.isr=1

log.flush.interval.messages=10000 
log.flush.interval.ms=1000

log.retention.hours=168 
log.retention.check.interval.ms=300000

zookeeper.connect=10.64.223.70:2181
zookeeper.connection.timeout.ms=18000

group.initial.rebalance.delay.ms=0 
max.connection.per.ip=100 
max.connections=100
listener.name.internal.max.connections=100
request.timeout.ms=180000
connections.max.idle.ms=300000

AdminAPI创建Topic代码

public void createTopics(Collection<NewTopic> newTopics, CreateTopicsOptions createTopicsOptions) {
    Properties props = new Properties();
    props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, brokers);
    AdminClient adminClient = KafkaAdminClient.create(props);
    adminClient.createTopics(newTopics, createTopicsOptions);
}

注:该方法在main方法中调用


问题排查与解决思路

1. AdminClient未关闭导致连接泄漏

观察AdminAPI代码,每次调用createTopics都会创建新的AdminClient但未调用close()方法。AdminClient内部会维护与Broker的TCP连接,不关闭的话这些连接会一直存在,导致连接数持续增长。

修复方案:使用完AdminClient后必须关闭:

public void createTopics(Collection<NewTopic> newTopics, CreateTopicsOptions createTopicsOptions) {
    Properties props = new Properties();
    props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, brokers);
    try (AdminClient adminClient = KafkaAdminClient.create(props)) {
        adminClient.createTopics(newTopics, createTopicsOptions);
        // 等待创建完成(可选,根据业务需求)
        adminClient.createTopics(newTopics, createTopicsOptions).all().get();
    } catch (InterruptedException | ExecutionException e) {
        // 处理异常
        logger.error("Failed to create topics", e);
        Thread.currentThread().interrupt();
    }
}

2. Log4j2 KafkaAppender连接配置问题

Log4j2的KafkaAppender如果配置不当,也可能创建多个Producer实例导致连接泄漏。检查KafkaAppender的配置,确保它使用的是单例Producer,或者配置了合适的连接回收策略。

建议配置:在log4j2.xml中确保KafkaAppender的producerConfig复用全局的Producer配置,并且设置合理的lingerMs、batchSize等参数,避免频繁创建连接。

3. Producer连接参数优化

虽然当前Producer是单例,但可以添加连接相关的配置来控制连接生命周期:

// 添加以下配置到Producer属性中
props.put(CommonClientConfigs.CONNECTIONS_MAX_IDLE_MS_CONFIG, 300000); // 5分钟闲置后关闭连接
props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5); // 控制每个连接的未响应请求数
props.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, 30000); // 缩短请求超时时间

4. 确认连接管理器单例正确性

检查KafkaConnectionManager的单例实现是否存在问题,确保整个应用中只创建了一个Producer实例。可以在日志中打印Producer的哈希值,确认每次获取的都是同一个实例。

5. Broker连接数配置调整

当前Broker配置中max.connection.per.ip=100,但应用连接数已经超过这个值,说明该配置可能未生效或者存在其他连接来源。可以适当调高该值,同时确保connections.max.idle.ms配置生效,让Broker主动关闭闲置连接。


内容的提问来源于stack exchange,提问作者Palash Gupta

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 16:01:01