Kafka消费异常求助:从uvtopic1-0获取记录时出错
错误信息
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)
问题分析
- 无限递归调用:
topicListener方法在finally块中递归调用自身,会导致栈溢出,且异常发生后会反复触发相同错误,无法正常恢复。 - 反序列化异常:该错误核心原因是某条消息无法被
StringDeserializer解析(比如消息不是合法UTF-8字符串、消息内容损坏),消费者无法读取该消息,进而抛出异常。 - 线程模型不规范: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

