Apache Storm KafkaSpout输出Tuple中value字段的声明位置说明
问题解答
StationParsingBolt中通过input.getStringByField("value")获取消息内容时用到的value字段,是Storm官方KafkaSpout组件默认定义的输出字段,不需要在当前项目的业务代码中手动声明。
- 你使用的
org.apache.storm.kafka.spout.KafkaSpout组件内置了默认的Tuple输出规则:从Kafka消费到单条消息后,会自动将消息封装为固定结构的Tuple,默认包含以下字段:topic:消息所属的Kafka主题名partition:消息所在的Kafka分区编号offset:消息在分区内的存储偏移量key:Kafka消息的键,对应生产者发送消息时传入的key参数value:Kafka消息的实际内容,即Python程序发送到Kafka的消息体
- 你在主类中构建
KafkaSpoutConfig时,没有通过setRecordTranslator方法自定义Tuple字段映射规则,因此KafkaSpout直接使用默认配置输出Tuple,下游Bolt可以直接通过value字段名获取Kafka的原始消息内容。 - 如果有自定义字段名、调整Tuple结构的需求,可以在构建KafkaSpoutConfig时传入自定义的消息转换逻辑实现,当前示例项目使用默认配置即可满足需求,不需要额外修改。
附:相关核心代码
主类拓扑定义代码
package velos; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.storm.Config; import org.apache.storm.LocalCluster; import org.apache.storm.StormSubmitter; import org.apache.storm.generated.AlreadyAliveException; import org.apache.storm.generated.AuthorizationException; import org.apache.storm.generated.InvalidTopologyException; import org.apache.storm.generated.StormTopology; import org.apache.storm.kafka.spout.KafkaSpout; import org.apache.storm.kafka.spout.KafkaSpoutConfig; import org.apache.storm.topology.TopologyBuilder; import org.apache.storm.topology.base.BaseWindowedBolt; import org.apache.storm.tuple.Fields; public class App { public static void main(String[] args) throws AlreadyAliveException, InvalidTopologyException, AuthorizationException { TopologyBuilder builder = new TopologyBuilder(); KafkaSpoutConfig.Builder<String, String> spoutConfigBuilder = KafkaSpoutConfig.builder("localhost:9092", "velib-stations"); spoutConfigBuilder.setProp(ConsumerConfig.GROUP_ID_CONFIG, "city-stats"); // 此处未自定义RecordTranslator,使用KafkaSpout默认输出字段 KafkaSpoutConfig<String, String> spoutConfig = spoutConfigBuilder.build(); builder.setSpout("stations", new KafkaSpout<String, String>(spoutConfig)); builder.setBolt("station-parsing", new StationParsingBolt()).shuffleGrouping("stations"); builder.setBolt("city-stats", new CityStatsBolt().withTumblingWindow(BaseWindowedBolt.Duration.of(1000 * 60 * 5))) .fieldsGrouping("station-parsing", new Fields("city")); builder.setBolt("save-results", new SaveResultsBolt()).fieldsGrouping("city-stats", new Fields("city")); StormTopology topology = builder.createTopology(); Config config = new Config(); config.setMessageTimeoutSecs(60 * 30); String topologyName = "Velos"; if (args.length > 0 && args[0].equals("remote")) { StormSubmitter.submitTopology(topologyName, config, topology); } else { LocalCluster cluster = new LocalCluster(); cluster.submitTopology(topologyName, config, topology); } } }
StationParsingBolt类代码
package velos; import java.util.Map; import org.apache.storm.task.OutputCollector; import org.apache.storm.task.TopologyContext; import org.apache.storm.topology.OutputFieldsDeclarer; import org.apache.storm.topology.base.BaseRichBolt; import org.apache.storm.tuple.Fields; import org.apache.storm.tuple.Tuple; import org.apache.storm.tuple.Values; import org.apache.storm.shade.org.json.simple.JSONObject; import org.apache.storm.shade.org.json.simple.parser.JSONParser; import org.apache.storm.shade.org.json.simple.parser.ParseException; public class StationParsingBolt extends BaseRichBolt { private OutputCollector outputCollector; @Override public void prepare(Map stormConf, TopologyContext context, OutputCollector collector) { outputCollector = collector; } @Override public void execute(Tuple input) { try { process(input); } catch (ParseException e) { e.printStackTrace(); outputCollector.fail(input); } } public void process(Tuple input) throws ParseException { JSONParser jsonParser = new JSONParser(); // 直接取KafkaSpout默认输出的value字段获取消息内容 JSONObject obj = (JSONObject)jsonParser.parse(input.getStringByField("value")); String contract = (String)obj.get("contract_name"); Long availableStands = (Long)obj.get("available_bike_stands"); Long stationNumber = (Long)obj.get("number"); outputCollector.emit(new Values(contract, stationNumber, availableStands)); outputCollector.ack(input); } @Override public void declareOutputFields(OutputFieldsDeclarer declarer) { declarer.declare(new Fields("city", "station_id", "available_stands")); } }
内容的提问来源于stack exchange,提问作者GRAN
相关产品推荐
相关产品推荐

