【问题标题】:2 spouts one Storm Bolt2个喷出1个风暴箭
【发布时间】:2016-03-07 02:13:42
【问题描述】:

我想创建一个包含 2 个具有 2 个不同主题的 kafkaSpout 的拓扑,并基于螺栓中的 sourceComponent 将这 2 个 spout 合并为一个流。

public class Topology {
private static final String topic1 = Real2
private static final String topic2 = Real1


public static void main(String[] args) throws AlreadyAliveException,
        InvalidTopologyException, IOException {

    BasicConfigurator.configure();

    String zookeeper_root = "";
    SpoutConfig kafkaConfig1 = new SpoutConfig(localhost:2181, topic1,
            zookeeper_root, "Real1KafkaSpout");

    SpoutConfig kafkaConfig2 = new SpoutConfig(localhost:2181, topic2,
            zookeeper_root, "Real2KafkaSpout");


    kafkaConfigRealTime.scheme = new SchemeAsMultiScheme(new StringScheme());
    kafkaConfigRealTime.forceFromStart = true;

    kafkaConfigHistorical.scheme = new SchemeAsMultiScheme(
            new StringScheme());
    kafkaConfigHistorical.forceFromStart = true;



    TopologyBuilder builder = new TopologyBuilder();

    builder.setSpout("Real1", new KafkaSpout(
            kafkaConfig1), 2);
    builder.setSpout("Real2", new KafkaSpout(
            kafkaConfig2), 2);

    builder.setBolt("StreamMerging", new StreamMergingBolt(), 2)
            .setNumTasks(2).shuffleGrouping("Real1")
            .shuffleGrouping("Real2");



    Config config = new Config();
    config.put("hdfs.config", yamlConf);

    config.setDebug(false);
    config.setMaxSpoutPending(10000);

    if (args.length == 0) {
        LocalCluster cluster = new LocalCluster();

        cluster.submitTopology("Topology", config,
                builder.createTopology());
        cluster.killTopology("Topology");
        cluster.shutdown();
    } else {

        StormSubmitter.submitTopology(args[0], config,
                builder.createTopology());
    }

    try {
        Thread.sleep(6000);
    } catch (InterruptedException ex) {
        ex.printStackTrace();
    }

}

}

在我正在做的螺栓执行方法中

public void execute(Tuple input, BasicOutputCollector collector) {
    String id = input.getSourceComponent();
    System.out.println("Stream Id in StreamMergingBolt is " + "---->" + id);

    }

所以我想将来自每个流的元组存储到单独的文件中 那就是我想将 Real1KafkaSpout 的元组存储到 file1 和 Real2KafkaSpout 到 file2 。我怎么能做到这一点我被打动了

【问题讨论】:

    标签: java apache-kafka apache-storm


    【解决方案1】:

    我通过以下代码执行此操作时的有线结果

    public void execute(Tuple input, BasicOutputCollector collector) {
    String id = input.getSourceComponent();
    if(id.equals("Real1")) {
       String json = input.getString(0);
      //writetoHDFS1(json)
       } else {
          assert (id.equals("Real2");
        String json = input.getString(0);
     //writetoHDFS2(json)
      }
     }
    

    【讨论】:

    • “有线结果”是什么意思?或者它是一个错字,你的意思是“奇怪的结果”? “以上代码在大多数情况下都有效”是什么意思?什么时候不起作用?为什么不呢?
    • @Matthias J. Sax 但是拓扑中有太多失败的元组......我怎样才能丢弃失败的元组......
    【解决方案2】:

    你可以这样做:

    public void execute(Tuple input, BasicOutputCollector collector) {
        String id = input.getSourceComponent();
        if(id.equals("Real1")) {
            // store into file1
        } else {
            assert (id.equals("Real2");
            // store in file2
        }
    }
    

    您将在Bolt.open(...) 中打开这两个文件。

    但是,我想知道您为什么要使用单个拓扑来执行此操作。如果您只将来自 Kafka 源 1 的数据写入文件 1,将来自 Kafka 源 2 的数据写入文件 2,您可以简单地创建两个拓扑......(当然,您只需对其进行一次编程,并针对两种情况进行不同的配置) .

    【讨论】:

    • 你是对的.. 但对于我的拓扑结构,来自 2 个不同来源的元组需要在螺栓中处理......但在处理之前,我需要在处理之后分别存储这些元组来自 2 个源的元组。数据会有所不同,但数据的结构是相同的。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2014-08-16
    • 2016-01-13
    • 1970-01-01
    • 1970-01-01
    • 2022-12-09
    • 2018-10-26
    • 1970-01-01
    相关资源
    最近更新 更多