Java应用使用Kafka Producer出现大量TCP-ESTABLISHED连接的原因排查
我在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

