Kafka消费者启动后消费少量消息即报Node Disconnected故障求助
Kafka消费者出现
nodes disconnected并停止消费的排查思路 问题描述
我有一个Kafka消费者应用,能够正常启动,但在消费少量消息后总会出现nodes disconnected提示并停止消费。我怀疑是不是因为消费者处理单条记录的耗时过长?但根据配置,我每次仅拉取5条记录,且单条记录的处理耗时并不久。
消费者配置代码
package net.abc.com; import lombok.extern.slf4j.Slf4j; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Profile; import org.springframework.kafka.annotation.EnableKafka; import org.springframework.kafka.core.reactive.ReactiveKafkaConsumerTemplate; import reactor.kafka.receiver.ReceiverOptions; import java.util.Collections; import java.util.HashMap; import java.util.Map; @Profile("default") @Configuration @EnableKafka @Slf4j public class KafkaConsumerConfigForTestTopics { private static final String SECURITY_PROTOCOL = "security.protocol"; private static final String SASL_MECHANISM = "sasl.mechanism"; private static final String SASL_JAAS_CONFIG = "sasl.jaas.config"; @Value("${kafka.bootstrap-servers}") private String bootstrapServers; @Value("${kafka.username:}") private String username; @Value("${kafka.password:}") private String kafkaSecretPass; @Value("${kafka.consumer-group}") private String consumerGroup; @Value("${kafka.login-module}") private String loginModule; @Value("${kafka.security-protocol}") private String securityProtocol; @Value("${kafka.sasl-mechanism:PLAIN}") private String saslMechanism; @Value("${kafka.listener.concurrency.count:1}") private int concurrencyCount; public Map<String, Object> consumerConfigs(String clientId) { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); props.put("key.deserializer", "org.apache.kafka.common.serialization.ByteArrayDeserializer"); props.put("value.deserializer", "org.apache.kafka.common.serialization.ByteArrayDeserializer"); props.put(ConsumerConfig.GROUP_ID_CONFIG, consumerGroup); props.put(SECURITY_PROTOCOL, securityProtocol); props.put(SASL_MECHANISM, saslMechanism); props.put(SASL_JAAS_CONFIG, String.format("%s required username=\"%s\" password=\"%s\" ;", loginModule, username, kafkaSecretPass)); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 5); props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 900000); props.put(ConsumerConfig.CLIENT_ID_CONFIG, clientId); props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG, "org.apache.kafka.clients.consumer.CooperativeStickyAssignor"); return props; } @Bean(name = "customer-group") public ReceiverOptions<byte[], byte[]> kafkaReceiverOptionsCustomerGroup(@Value(value = "${kafka.consumer.topic-customer-group}") String topic, KafkaProperties kafkaProperties) { ReceiverOptions<byte[], byte[]> basicReceiverOptions = ReceiverOptions.create(consumerConfigs("customer-group-consumer")); return basicReceiverOptions.subscription(Collections.singletonList(topic)); } @Bean(name = "customer-group-template") public ReactiveKafkaConsumerTemplate<byte[], byte[]> reactiveKafkaConsumerTemplateCustomerGroup(@Qualifier("customer-group") ReceiverOptions<byte[], byte[]> kafkaReceiverOptions) { return new ReactiveKafkaConsumerTemplate<>(kafkaReceiverOptions); } // Like this I have four listeners in this project... }
错误日志
[Consumer clientId=customer-group-consumer, groupId=abc.cdc.xyz.consumerGroup.v1] Node 25 disconnected. [Consumer clientId=customer-group-consumer, groupId=abc.cdc.xyz.consumerGroup.v1] Error sending fetch request (sessionId=609702141, epoch=687281) to node 56: [Consumer clientId=customer-group-consumer, groupId=abc.cdc.xyz.consumerGroup.v1] Cancelled in-flight FETCH request with correlation id 1268856 due to node 56 being disconnected (elapsed time since creation: 1ms, elapsed time since send: 1ms, request timeout: 30000ms)
可能的原因及解决方法
- 网络连通性问题:错误日志明确显示节点断开、请求发送失败,优先排查消费者与Kafka broker之间的网络。可以用
telnet <broker-ip> <port>或nc -zv <broker-ip> <port>持续测试连通性,同时查看broker端日志,确认是否有节点过载、连接超时的记录。另外,检查防火墙是否限制了消费者与broker之间的通信。 - 心跳与会话超时配置不合理:当前配置只设置了
MAX_POLL_INTERVAL_MS_CONFIG,但Reactive Kafka底层依赖的session.timeout.ms(默认30秒)和heartbeat.interval.ms(默认3秒)如果配置过小,可能导致broker认为消费者失联。建议调整:props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 300000); // 5分钟 props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 10000); // 10秒 - Reactive流处理阻塞:如果消息处理的Reactive流中存在
block()等阻塞操作,会占用消费者线程,导致无法正常发送心跳。检查消息处理逻辑,确保所有操作都是非阻塞的Reactive风格。 - SASL认证会话问题:使用SASL认证时,如果凭证过期或认证失败,broker会断开连接。检查用户名密码是否正确,查看broker的认证日志(如
kafka-authorizer.log)是否有认证失败记录。部分SASL机制支持自动重新认证,确保客户端配置了相关参数。 - 版本兼容性问题:如果Kafka客户端版本与broker版本差距过大,可能存在协议不兼容。比如客户端用了2.x版本而broker是0.10.x,或者反过来。尽量使用与broker相同大版本的客户端。
内容的提问来源于stack exchange,提问作者Sree B
相关产品推荐
相关产品推荐

