Java自学者论坛

 找回密码
 立即注册

手机号码,快捷登录

恭喜Java自学者论坛(https://www.javazxz.com)已经为数万Java学习者服务超过8年了!积累会员资料超过10000G+
成为本站VIP会员,下载本站10000G+会员资源,会员资料板块,购买链接:点击进入购买VIP会员

JAVA高级面试进阶训练营视频教程

Java架构师系统进阶VIP课程

分布式高可用全栈开发微服务教程Go语言视频零基础入门到精通Java架构师3期(课件+源码)
Java开发全终端实战租房项目视频教程SpringBoot2.X入门到高级使用教程大数据培训第六期全套视频教程深度学习(CNN RNN GAN)算法原理Java亿级流量电商系统视频教程
互联网架构师视频教程年薪50万Spark2.0从入门到精通年薪50万!人工智能学习路线教程年薪50万大数据入门到精通学习路线年薪50万机器学习入门到精通教程
仿小米商城类app和小程序视频教程深度学习数据分析基础到实战最新黑马javaEE2.1就业课程从 0到JVM实战高手教程MySQL入门到精通教程
查看: 736|回复: 0

kafka使用getOffsetsBefore()获取获取offset异常分析

[复制链接]
  • TA的每日心情
    奋斗
    2024-4-6 11:05
  • 签到天数: 748 天

    [LV.9]以坛为家II

    2034

    主题

    2092

    帖子

    70万

    积分

    管理员

    Rank: 9Rank: 9Rank: 9

    积分
    705612
    发表于 2021-5-28 16:40:09 | 显示全部楼层 |阅读模式

    根据时间戳获取kafka的topic的偏移量,结果获取的偏移量量数据组的长度为0,就会出现如下的数组下标越界的异常,实现的原理是使用了kafka的getOffsetsBefore()方法:

     

    Exception in thread "main" java.lang.ArrayIndexOutOfBoundsException : 0

         at co.gridport.kafka.hadoop.KafkaInputFetcher.getOffset(KafkaInputFetcher.java:126)

         at co.gridport.kafka.hadoop.TestKafkaInputFetcher.testGetOffset(TestKafkaInputFetcher.java:68)

         at co.gridport.kafka.hadoop.TestKafkaInputFetcher.main(TestKafkaInputFetcher.java:80)

    OffsetResponse(0,Map([barrage_detail,0] -> error: kafka.common.UnknownException offsets: ))

    源码如下:

    /*

          * 得到partition的offset Finding Starting Offset for Reads

          */

    public Long getOffset(Long time) throws IOException {

               TopicAndPartition topicAndPartition = new TopicAndPartition(this.topic , this.partition );

               Map<TopicAndPartition, PartitionOffsetRequestInfo> requestInfo = new HashMap<TopicAndPartition, PartitionOffsetRequestInfo>();

    requestInfo.put( topicAndPartition, new PartitionOffsetRequestInfo(time, 1));

               kafka.javaapi.OffsetRequest request = new kafka.javaapi.OffsetRequest(

    requestInfo, kafka.api.OffsetRequest.CurrentVersion(), this. client_id);

               OffsetResponse response = this. kafka_consumer.getOffsetsBefore( request);

    if ( response.hasError()) {

    log.error( "Error fetching data Offset Data the Broker. Reason: " + response.errorCode(this.topic, this. partition));

    throw new IOException ( "Error fetching kafka Offset by time:" + response.errorCode(this.topic, this. partition));

               }

    //         if (response.offsets(this.topic, this.partition).length == 0){

    //              return getOffset(kafka.api.OffsetRequest

    //                         .EarliestTime());

    //         }

    return response.offsets( this. topic, this. partition)[0];

         }

    返回的response对象会有error: kafka.common.UnknownException offsets如下异常:

    OffsetResponse(0,Map([barrage_detail,0] -> error: kafka.common.UnknownException offsets: ))

    同时呢,response.hasError()检查不到error。

    是什么原因造成了response.offsets(this.topic,this.partition)的返回数组的长度为0呢?

    分析了getOffsetsBefore()方法的源码,并做源码了大量的测试,终于重现了这种情况?

    1.getOffsetsBefore()的功能以及实现原理:

    getOffsetsBefore的功能是返回某个时间点前的maxOffsetNum个offset(时间点指的是kafka日志文件的最后修改时间,offset指的是log文件名中的offset,这个offset是该日志文件的第一条记录的offset,即base offset;maxNumOffsets参数即返回结果的最大个数,如果该参数为2,就返回指定时间点前的2个offset,如果是负数,就报逻辑错误,怎么能声明一个长度为负数的数组呢,呵呵);

    根据这个实现原理,所以返回的结果长度为0是合理的,反映的是这个时间点前没有kafka日志这种情况,该情况自然就没有offset了。

    说明我们指定的时间参数太早了,正常的时间范围为:最早的时间之后的时间参数都可以有返回值。

    其实合理的处理方式应该为如果这个时间点前没有值,就返回最早的offset了,对api的使用者就友好多了我们可以自己来实现这个功能。

    代码如下:

     

    /*

          * 得到partition的offset Finding Starting Offset for Reads

          */

    public Long getOffset(Long time ) throws IOException {

               TopicAndPartition topicAndPartition = new TopicAndPartition(this .topic , this.partition);

               Map<TopicAndPartition, PartitionOffsetRequestInfo> requestInfo = new HashMap<TopicAndPartition, PartitionOffsetRequestInfo>();

    requestInfo.put( topicAndPartition, new PartitionOffsetRequestInfo(time, 1));

               kafka.javaapi.OffsetRequest request = new kafka.javaapi.OffsetRequest(

    requestInfo, kafka.api.OffsetRequest.CurrentVersion(), this. client_id);

               OffsetResponse response = this. kafka_consumer.getOffsetsBefore( request);

    if ( response.hasError()) {

    log.error( "Error fetching data Offset Data the Broker. Reason: " + response.errorCode(this.topic, this. partition));

    throw new IOException ( "Error fetching kafka Offset by time:" + response.errorCode(this.topic, this. partition));

               }

                //如果返回的数据长度为0,就获取最早的offset。

    if ( response.offsets( this. topic, this. partition). length == 0){

    return getOffset(kafka.api.OffsetRequest

                               . EarliestTime());

               }

    return response.offsets( this. topic, this. partition)[0];

         }

    哎...今天够累的,签到来了1...
    回复

    使用道具 举报

    您需要登录后才可以回帖 登录 | 立即注册

    本版积分规则

    QQ|手机版|小黑屋|Java自学者论坛 ( 声明:本站文章及资料整理自互联网,用于Java自学者交流学习使用,对资料版权不负任何法律责任,若有侵权请及时联系客服屏蔽删除 )

    GMT+8, 2024-5-13 11:34 , Processed in 0.068228 second(s), 29 queries .

    Powered by Discuz! X3.4

    Copyright © 2001-2021, Tencent Cloud.

    快速回复 返回顶部 返回列表