【问题标题】:Field Grouping for a Kafka SpoutKafka Spout 的字段分组
【发布时间】:2014-11-20 09:43:59
【问题描述】:

可以对 kafka spout 发出的元组进行字段分组吗?如果是,Storm 是如何知道 Kafka 记录中的字段的?

【问题讨论】:

    标签: apache-storm apache-kafka


    【解决方案1】:

    Storm 中的字段分组(和一般分组)是针对螺栓的,而不是针对喷口的。这是通过InputDeclarer 类完成的。
    当您在TopologyBuilder 上调用setBolt() 时,将返回InputDeclarer

    【讨论】:

    • 我的错,我的意思是螺栓。那就是我有一个卡夫卡喷口,它将把元组发送到后续的螺栓。现在对于风暴分布中包含的 kafka spout,我必须首先知道它发出的字段。这些字段 id 是否与 kafka 发布者发布的相同?
    【解决方案2】:

    Kafka Spout 像任何其他组件一样声明其输出字段。我的解释是基于KafkaSpout当前的implementation

    在 KafkaSpout.java 类中,我们看到 declareOutputFields 方法调用 KafkaConfig Scheme 的 getOutputFields() 方法。

    @Override
    public void declareOutputFields(OutputFieldsDeclarer declarer) {
        declarer.declare(_spoutConfig.scheme.getOutputFields());
    }
    

    KafkaConfig默认使用RawMultiScheme就是这样实现这个方法的。

      @Override
      public Fields getOutputFields() {
        return new Fields("bytes");
      }
    

    那么这是什么意思呢?如果你用 fieldGrouping 声明了从 KafkaSpout 读取元组的 bolt,你就知道 每个包含 equals 字段“bytes”的元组都将由同一个任务执行。如果你想发出任何字段,你应该根据你的需要实施新的方案。

    【讨论】:

      【解决方案3】:

      TL:DR KafkaSpout 的默认实现在declareOutputFields 中声明了以下输出字段:

      new Fields("topic", "partition", "offset", "key", "value");

      所以在构建拓扑代码时直接做:

      topologyBuilder.setSpout(spoutName, mySpout, parallelismHintSpout);
      topologyBuilder.setBolt(boltName, myBolt, parallelismHintBolt).fieldsGrouping(spoutName, new Fields("key"));
      

      详细信息:稍微研究一下代码就会发现:

      在 Kafka Spout 中,declareOutputFieldsimplemented,如下所示:

      @Override
      public void declareOutputFields(OutputFieldsDeclarer declarer) {
          RecordTranslator<K, V> translator = kafkaSpoutConfig.getTranslator();
          for (String stream : translator.streams()) {
              declarer.declareStream(stream, translator.getFieldsFor(stream));
          }
      }
      

      它从RecordTranslator 接口获取字段,并从kafkaSpoutConfig 获取其实例,即KafkaSpoutConfig&lt;K, V&gt;KafkaSpoutConfig&lt;K, V&gt;CommonKafkaSpoutConfig 扩展而来(这在 1.1.1 版本中略有不同)。这个returnsDefaultRecordTranslator的建造者。如果你检查这个类中的字段implementation,你会发现:

      public static final Fields FIELDS = new Fields("topic", "partition", "offset", "key", "value");
      

      所以我们可以直接在拓扑代码中的字段分组中使用Fields("key")

      topologyBuilder.setBolt(boltName, myBolt, parallelismHintBolt).fieldsGrouping(spoutName, new Fields("key"));
      

      【讨论】:

        猜你喜欢
        • 2014-12-03
        • 2013-06-24
        • 2015-03-13
        • 2017-08-24
        • 2020-05-01
        • 2013-08-18
        • 1970-01-01
        • 1970-01-01
        • 2015-12-30
        相关资源
        最近更新 更多