【问题标题】:Creating a KafkaSpout with a consumer group id使用消费者组 ID 创建 KafkaSpout
【发布时间】:2017-02-26 12:45:26
【问题描述】:

一段时间以来,我一直在尝试了解如何创建一组订阅单个主题的 kafkaspout。我发现一些引用 spoutconfig 中的 id 的来源可以用作组 id。我的主要问题是我如何知道创建的 spout 是否作为一个组起作用。我还想知道setSpout() 中的 paralellism_hint 是否是创建一组 paralellism_hint 数量的 spout 的人。请赐教。

我的代码是这样的

private KafkaSpout buildKafkaSpout(String zkTopic, String zkRoot, String groupId) {

        String zkConnString = "localhost:2181";
        BrokerHosts hosts = new ZkHosts(zkConnString);
        SpoutConfig kafkaSpoutConfig = new SpoutConfig (hosts, zkTopic,zkRoot, groupId );
        kafkaSpoutConfig.scheme = new SchemeAsMultiScheme(new StringScheme());
        KafkaSpout kafkaSpout = new KafkaSpout(kafkaSpoutConfig);
        return kafkaSpout;
    }

private void buildTopologyForGroupOne(TopologyBuilder builder ) {

        String zkTopic = "topic1";
        String zkRoot = "/topic1";
        String groupId = "group1";
        List<String> zkSpoutIds = new ArrayList<String>();
        zkSpoutIds.add("word_count-spout");
        zkSpoutIds.add("total_word_count-spout");

        for(String spoutId:zkSpoutIds ){
            KafkaSpout kafkaSpout = buildKafkaSpout(zkTopic, zkRoot, groupId);
            builder.setSpout(spoutId.concat("_"+groupId), kafkaSpout,2);
        }
        builder.setBolt("word_split-bolt_group_1",new SplitBolt()).shuffleGrouping("word_count-spout"+"_"+groupId);
        builder.setBolt("split_count-bolt_group_1",new CountBolt()).shuffleGrouping("word_split-bolt_group_1");
        builder.setBolt("total_word_count-bolt_group_1",new TotalWordCountBolt()).shuffleGrouping("total_word_count-spout"+"_"+groupId);

    }

现在我如何知道我创建的两个 spout(word_count-spout,total_word_count-spout) 是否作为一个组。作为一个组,我的意思是如果创建了一个新的 spout,zookeeper 将重新排列分区。

提前致谢

【问题讨论】:

  • 终于找到了来源。可以使用新版本的storm apache-storm-1.0.2 创建一个带有consumer id 的KafkaSpout。

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


【解决方案1】:

storm-kafka 使用 SimpleConsumer 而不是高级消费者 API,因此 group.id 并不重要,因为它是用于使用此设置实现高可用性和故障转移的高级消费者。

回到你的问题,我更倾向于认为你的代码在运行时会发生错误,因为两个消费者实例将争夺同一个 zk 路径锁。

【讨论】:

  • 我的意图是为 kafkaSpout 使用组 ID。您的意思是说这不可能吗?
  • 出于性能考虑,可以按照Storm的方式设置worker数/executor数/task数,以获得更好的吞吐量。设置 group.id 对 KafkaSpout 没有影响。除此之外应该是客户端 ID,它告诉 ZK 在哪里存储偏移量,而不是你认为的组 ID。
  • 感谢您的澄清。我的目标是使一个喷口输出被多个螺栓消耗,这是可以做到的。但是,这种实现中的问题是,如果一个螺栓无法确认元组,元组也由喷口重播到其他每个螺栓。我想避免这种情况。如果你能建议我一条出路,我将非常感激。我一直在尝试使用流名称来寻找解决方案。还没有祝你好运。
  • 我记得 Storm 只重放存在无法确认的失败元组的流。您可能会声明多个流以使不同的螺栓订阅不同的流。这样做可能会避免您担心的情况。
  • 我一直在尝试相同的方法。但问题是即使螺栓确认,喷口无法确认并且消息不断重播,因此我无法测试一个螺栓是否失败会发生什么。而且我也无法创建一个读取的喷口像KafkaSpout一样从主题中获得帮助。请在这个方向上提供任何帮助。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2017-05-13
  • 1970-01-01
  • 2021-05-14
  • 1970-01-01
  • 2015-09-08
  • 2017-08-24
  • 1970-01-01
相关资源
最近更新 更多