You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.30 22:18:03