【发布时间】:2014-11-20 09:43:59
【问题描述】:
可以对 kafka spout 发出的元组进行字段分组吗?如果是,Storm 是如何知道 Kafka 记录中的字段的?
【问题讨论】:
可以对 kafka spout 发出的元组进行字段分组吗?如果是,Storm 是如何知道 Kafka 记录中的字段的?
【问题讨论】:
Storm 中的字段分组(和一般分组)是针对螺栓的,而不是针对喷口的。这是通过InputDeclarer 类完成的。
当您在TopologyBuilder 上调用setBolt() 时,将返回InputDeclarer。
【讨论】:
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”的元组都将由同一个任务执行。如果你想发出任何字段,你应该根据你的需要实施新的方案。
【讨论】:
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 中,declareOutputFields 是 implemented,如下所示:
@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<K, V>。 KafkaSpoutConfig<K, V> 从 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"));
【讨论】: