【问题标题】:Consumer fetch data returns OFFSET_OUT_OF_RANGE消费者获取数据返回 OFFSET_OUT_OF_RANGE
【发布时间】:2020-07-01 11:49:08
【问题描述】:

我有一个包含 3 个 kafka 代理的集群,主题名为 fallback_topic 只有一个consumerGroup从这个topic消费,并且这个consumerGroup中只有一个consumer

注入几条消息后,我可以看到消息已发布到 Kafka。 LogSize 已被新消息移动;但是,Consumer Offset 保持不变,并且不会消费任何消息。

下面是consumer.poll(3000) 运行时的日志。分区 (4, 7, 10) 收到了来自生产者的新消息,但是当消费者尝试读取它时,它报告了error=OFFSET_OUT_OF_RANGE

04:20:41.311 [kafka-coordinator-heartbeat-thread | uniqueConsumerGroup] DEBUG o.a.k.clients.FetchSessionHandler - [Consumer clientId=consumer-1, groupId=uniqueConsumerGroup] Node 654000 sent a full fetch response that created a new incremental fetch session 685508830 with 7 response partition(s)
04:20:41.311 [kafka-coordinator-heartbeat-thread | uniqueConsumerGroup] DEBUG o.a.k.c.consumer.internals.Fetcher - [Consumer clientId=consumer-1, groupId=uniqueConsumerGroup] Fetch READ_UNCOMMITTED at offset 1062 for partition fallback_topic-1 returned fetch data (error=NONE, highWaterMark=1062, lastStableOffset = -1, logStartOffset = 1062, abortedTransactions = null, recordsSizeInBytes=0)
04:20:41.311 [kafka-coordinator-heartbeat-thread | uniqueConsumerGroup] DEBUG o.a.k.c.consumer.internals.Fetcher - [Consumer clientId=consumer-1, groupId=uniqueConsumerGroup] Fetch READ_UNCOMMITTED at offset 124094 for partition fallback_topic-4 returned fetch data (error=OFFSET_OUT_OF_RANGE, highWaterMark=-1, lastStableOffset = -1, logStartOffset = -1, abortedTransactions = null, recordsSizeInBytes=0)
04:20:41.311 [kafka-coordinator-heartbeat-thread | uniqueConsumerGroup] DEBUG o.a.k.c.consumer.internals.Fetcher - [Consumer clientId=consumer-1, groupId=uniqueConsumerGroup] Fetch READ_UNCOMMITTED at offset 762 for partition fallback_topic-7 returned fetch data (error=OFFSET_OUT_OF_RANGE, highWaterMark=-1, lastStableOffset = -1, logStartOffset = -1, abortedTransactions = null, recordsSizeInBytes=0)
04:20:41.311 [kafka-coordinator-heartbeat-thread | uniqueConsumerGroup] DEBUG o.a.k.c.consumer.internals.Fetcher - [Consumer clientId=consumer-1, groupId=uniqueConsumerGroup] Fetch READ_UNCOMMITTED at offset 897 for partition fallback_topic-10 returned fetch data (error=OFFSET_OUT_OF_RANGE, highWaterMark=-1, lastStableOffset = -1, logStartOffset = -1, abortedTransactions = null, recordsSizeInBytes=0)

我的理解是当分区的leader改变了offset,而follower没有,那就是这个错误发生的时候。但是没有代理中断,所以消费者一直在使用同一个领导者。谁能帮我解释为什么会出现 OFFSET_OUT_OF_RANGE 错误。非常感谢。下面是我的代码,我跳过了consumer.commitAsync(),因为我的问题发生在提交之前。

        List<Event> events = new ArrayList<Event>();
        consumer.subscribe(Arrays.asList("fallback_topic"));
        ConsumerRecords<String, byte[]> records;
        
        do {
            logger.info("Start polling messages from " + topic);
            records = consumer.poll(3000);

            logger.info("done polling.");
            records.partitions().forEach(tp -> logger.info("found records from "+tp.topic()+"-"+tp.partition()));
            for (ConsumerRecord<String, byte[]> record : records) {
                Event event = EventKafkaSerializer.serializer.deserializeEvent(new ByteArrayInputStream(record.value()));
                logger.info(event.getId()+" "+event.getData().toString());
                events.add(event);
            }
           
        } while(records.count()>0);
        
        logger.info("Found total events "+events.size());

【问题讨论】:

    标签: apache-kafka kafka-consumer-api kafka-producer-api


    【解决方案1】:

    找出原因。

    最后忘记运行consumer.close()

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2020-10-22
      • 1970-01-01
      • 1970-01-01
      • 2020-12-06
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多