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入门到精通教程
查看: 574|回复: 0

KafkaSpout 重复消费问题解决

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

    [LV.9]以坛为家II

    2034

    主题

    2092

    帖子

    70万

    积分

    管理员

    Rank: 9Rank: 9Rank: 9

    积分
    705612
    发表于 2021-4-20 11:40:36 | 显示全部楼层 |阅读模式

    使用https://github.com/nathanmarz/storm-contrib来对接Kafka0.7.2时, 发现kafkaSpout总会进行数据重读, 配置都无问题, 也没报错

    进行debug之后, 发现是由于自己写的blot继承于IBolt, 但自己没有在代码中显示的调用collector.ack(); 导致kafkaSpout一直认为emitted的数据有问题, 超时之后进行数据重发

    KafkaSpout中关键代码如下:

    PartitionManager.java

    public void commit() {
    LOG.info("Committing offset for " + _partition);
    long committedTo;
    if(_pending.isEmpty()) {
    committedTo = _emittedToOffset;
    } else {
    committedTo = _pending.first();
    }
    if(committedTo!=_committedTo) {
    LOG.info("Writing committed offset to ZK: " + committedTo);
    
    Map<Object, Object> data = (Map<Object,Object>)ImmutableMap.builder()
    .put("topology", ImmutableMap.of("id", _topologyInstanceId,
    "name", _stormConf.get(Config.TOPOLOGY_NAME)))
    .put("offset", committedTo)
    .put("partition", _partition.partition)
    .put("broker", ImmutableMap.of("host", _partition.host.host,
    "port", _partition.host.port))
    .put("topic", _spoutConfig.topic).build();
    _state.writeJSON(committedPath(), data);
    
    LOG.info("Wrote committed offset to ZK: " + committedTo);
    _committedTo = committedTo;
    }
    LOG.info("Committed offset " + committedTo + " for " + _partition);
    }

    如果Bolt不进行ack, 则红色代码处的offsetNumber永远相等, 导致一直不进行offset的回写操作

     

    解决方案:

    1. IBolt中显式调用collector.ack();

    2. 使用帮你封装好的BaseBasicBlot, 它会帮你自动调用ack的

    关于Ack的问题, 可以参考我的翻译和官网文章: http://www.cnblogs.com/zhwbqd/p/3960991.html

     

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

    使用道具 举报

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

    本版积分规则

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

    GMT+8, 2024-5-18 07:18 , Processed in 0.062043 second(s), 29 queries .

    Powered by Discuz! X3.4

    Copyright © 2001-2021, Tencent Cloud.

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