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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 09:02:53