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

Kafka消费异常求助:从uvtopic1-0获取记录时出错

Kafka消费者异常排查与解决

错误信息

Received exception when fetching the next record from uvtopic1-0. If needed, please seek past the record to continue consumption.

消费者实现代码

public class KafkaConsumerPoc2 {

    private static final Integer count = 10000;

    public static final String KAFKA_SERVER_URL = "0.0.0.0";
    public static final int KAFKA_SERVER_PORT = 29092;

    static ExecutorService executorService;

    static Map<String, KafkaProducer<String, String>> topicWiseProducerMap = new HashMap<String, KafkaProducer<String, String>>();

    private static Map<String, KafkaConsumer<String, String>> topicConsumerMap = new HashMap<String, KafkaConsumer<String, String>>();

    private static Properties getConsumerProp(String topicName) {
        Properties props = new Properties();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, KAFKA_SERVER_URL + ":" + KAFKA_SERVER_PORT);
        String grupId = topicName + "-grup";
        props.put(ConsumerConfig.GROUP_ID_CONFIG, grupId);
        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "true");
        props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, "1000");
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());

        return props;
    }

    public static KafkaConsumer<String, String> getConsumerByTopicName(String topicName) {
        KafkaConsumer<String, String> consumer = topicConsumerMap.get(topicName);
        if (consumer == null) {
            consumer = registerConsumer(topicName);
        }
        return consumer;
    }

    public static KafkaConsumer<String, String> registerConsumer(String topicName) {
        Properties pro = getConsumerProp(topicName);
        KafkaConsumer<String, String> consumer = new KafkaConsumer<String, String>(pro);
        consumer.subscribe(Collections.singleton(topicName));
        topicConsumerMap.put(topicName, consumer);
        return consumer;
    }

    private static void startThreadForTopicListening(String topic) {
        executorService.submit(new Callable<Boolean>() {

            @Override
            public Boolean call() throws Exception {

                KafkaConsumer<String, String> consumer = getConsumerByTopicName(topic);

                topicListener(topic, consumer);
                return true;
            }
        });
    }

    public static void topicListener(String topic, KafkaConsumer<String, String> consumer) {
        try {

            System.out.println("************* Read message starts *****************************");

            ConsumerRecords<String, String> consumerRecords = consumer.poll(Duration.ofMillis(1000));

            for (ConsumerRecord<String, String> record : consumerRecords) {

                if (record.value() != null) {
                    System.out.println("Received message: (" + record.value() + ") at offset " + record.offset()
                            + " topic : " + record.topic());
                }
            }

            System.out.println("************* Read message ends *****************************");

        } catch (Exception e) {
            e.printStackTrace();
        } finally {
            topicListener(topic, consumer);
        }
    }

    public static void main(String[] args) {

        System.out.println("Main starts");

        Integer sizeOfExecutors = 2;

        executorService = Executors.newFixedThreadPool(sizeOfExecutors);

        startThreadForWriting(KafkaClientPoc2.topic1);
        startThreadForTopicListening(KafkaClientPoc2.topic1);

    }

    private static void startThreadForWriting(String topic) {
        executorService.submit(new Callable<Boolean>() {

            @Override
            public Boolean call() throws Exception {
                for (int i = 1; i <= count; i++) {
                    KafkaClientPoc2.writeSingleMsgInTopic(KafkaClientPoc2.topic1, "Msg:" + i);
                }
                return true;
            }
        });
    }

    public static void writeSingleMsgInTopic(String topicName, String msg) {
        System.out.println("#################### Write Msg starts ############################");
        KafkaProducer<String, String> producer = getProducer(topicName);
        try {
            ProducerRecord<String, String> record = new ProducerRecord<String, String>(topicName, msg);
            producer.send(record);
            producer.flush();
            System.out.println("writer > Sent message: (" + msg + ")");
        } catch (Exception e) {
            e.printStackTrace();
        }

        System.out.println("#################### Write Msg ends ############################");
    }

    public static KafkaProducer<String, String> getProducer(String topicName) {
        KafkaProducer<String, String> producer = topicWiseProducerMap.get(topicName);
        if (producer == null) {
            producer = registerProducer(topicName);
        }
        return producer;
    }

    public static Properties getProducerProp() {
        Properties prop = new Properties();
        prop.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "0.0.0.0:29092");
        prop.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        prop.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        prop.put(ProducerConfig.BATCH_SIZE_CONFIG, 1);

        return prop;
    }

    private static KafkaProducer<String, String> registerProducer(String topicName) {
        System.out.println("Creating new Producer");
        KafkaProducer<String, String> producer = new KafkaProducer<String, String>(getProducerProp());
        topicWiseProducerMap.put(topicName, producer);
        return producer;
    }

}

堆栈跟踪信息

org.apache.kafka.common.KafkaException: Received exception when fetching the next record from uvtopic1-0. If needed, please seek past the record to continue consumption.
at org.apache.kafka.clients.consumer.internals.Fetcher$CompletedFetch.fetchRecords(Fetcher.java:1598)
at org.apache.kafka.clients.consumer.internals.Fetcher$CompletedFetch.access$1700(Fetcher.java:1453)
at org.apache.kafka.clients.consumer.internals.Fetcher.fetchRecords(Fetcher.java:686)
at org.apache.kafka.clients.consumer.internals.Fetcher.fetchedRecords(Fetcher.java:637)
at org.apache.kafka.clients.consumer.KafkaConsumer.pollForFetches(KafkaConsumer.java:1276)
at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1237)
at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1210)
at com.vp.loaddata.vploaddata.poc2.KafkaConsumerPoc2.topicListener(KafkaConsumerPoc2.java:80)
at com.vp.loaddata.vploaddata.poc2.KafkaConsumerPoc2.topicListener(KafkaConsumerPoc2.java:101)


问题分析

  1. 无限递归调用:topicListener方法在finally块中递归调用自身,会导致栈溢出,且异常发生后会反复触发相同错误,无法正常恢复。
  2. 反序列化异常:该错误核心原因是某条消息无法被StringDeserializer解析(比如消息不是合法UTF-8字符串、消息内容损坏),消费者无法读取该消息,进而抛出异常。
  3. 线程模型不规范:Kafka消费者应在循环中持续调用poll,而非递归调用方法,递归会破坏消费者的状态管理逻辑。

解决方案

1. 替换递归为循环调用

修改topicListener方法,用while(true)循环替代递归,保证消费者持续运行且避免栈溢出:

public static void topicListener(String topic, KafkaConsumer<String, String> consumer) {
    // 用循环替代递归
    while (!Thread.currentThread().isInterrupted()) {
        try {
            System.out.println("************* Read message starts *****************************");
            ConsumerRecords<String, String> consumerRecords = consumer.poll(Duration.ofMillis(1000));
            
            for (ConsumerRecord<String, String> record : consumerRecords) {
                if (record.value() != null) {
                    System.out.println("Received message: (" + record.value() + ") at offset " + record.offset()
                            + " topic : " + record.topic());
                }
            }
            System.out.println("************* Read message ends *****************************");
            
        } catch (SerializationException e) {
            // 处理反序列化异常,跳过损坏的消息
            handleDeserializationError(consumer, e);
        } catch (Exception e) {
            e.printStackTrace();
            // 添加延迟避免频繁报错
            try {
                Thread.sleep(1000);
            } catch (InterruptedException ie) {
                Thread.currentThread().interrupt();
            }
        }
    }
}

2. 处理反序列化异常,跳过损坏消息

添加handleDeserializationError方法,手动将消费者偏移量移动到损坏消息的下一个位置,恢复消费:

private static void handleDeserializationError(KafkaConsumer<String, String> consumer, SerializationException e) {
    // 获取异常中的主题分区和偏移量信息
    String errorMsg = e.getMessage();
    if (errorMsg.contains("offset")) {
        // 解析主题、分区和偏移量(根据实际异常信息调整逻辑)
        String[] parts = errorMsg.split(" ");
        String topicPartitionStr = parts[parts.length - 3];
        long offset = Long.parseLong(parts[parts.length - 1]);
        
        String[] topicPartition = topicPartitionStr.split("-");
        String topic = topicPartition[0];
        int partition = Integer.parseInt(topicPartition[1]);
        
        // 移动到下一个偏移量
        TopicPartition tp = new TopicPartition(topic, partition);
        consumer.seek(tp, offset + 1);
        System.out.println("跳过损坏消息,偏移量移动到: " + (offset + 1));
    }
}

3. 优化消费者生命周期管理

在main方法中添加关闭钩子,确保程序退出时正确关闭消费者和线程池:

public static void main(String[] args) {
    System.out.println("Main starts");
    Integer sizeOfExecutors = 2;
    executorService = Executors.newFixedThreadPool(sizeOfExecutors);
    
    startThreadForWriting(KafkaClientPoc2.topic1);
    startThreadForTopicListening(KafkaClientPoc2.topic1);
    
    // 添加关闭钩子
    Runtime.getRuntime().addShutdownHook(new Thread(() -> {
        executorService.shutdown();
        try {
            if (!executorService.awaitTermination(5, TimeUnit.SECONDS)) {
                executorService.shutdownNow();
            }
        } catch (InterruptedException ie) {
            executorService.shutdownNow();
        }
        // 关闭所有消费者
        for (KafkaConsumer<String, String> consumer : topicConsumerMap.values()) {
            consumer.close();
        }
        System.out.println("资源已释放");
    }));
}

内容的提问来源于stack exchange,提问作者Urvesh Rathod

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 08:30:47