【问题标题】:join two streams based on common field in storm bolt基于风暴螺栓中的公共字段连接两个流
【发布时间】:2016-04-20 11:12:42
【问题描述】:

问题陈述: 我想加入来自两个不同 Kafka Spout(比如 S1 和 S2)的两个流,并希望根据其中的一些公共字段加入来自每个流的元组。 如果“S1”接收低于 json 作为元组

{"l7ProtocolID":"dhcp",
"packets_out":1,
"bytes_out":400,
"start_time":1454281199898,
"flow_sample":0,
"duration":102,
"path":["base","ip","udp","dhcp"],
"bytes_in":1200,
"l4":[{"client":"68","server":"67","level":0}],
"l2":[{"client":"52:54:00:50:04:B2","server":"FF:FF:FF:FF:FF:FF","level":0}],
"l3":[{"client":"::ffff:0.0.0.0","server":"::ffff:255.255.255.255","level":0}],
"flow_id":"81454281200000731489",
"applicationID":"dhcp",
"packets_in":1}

并且“S2”在 JSON 下作为元组接收

{"portGroupName":"dhcp",
"hypervisorName":1,
"bytes_out":400,
"monitoredIP":1454281199898,
"monitoredInstance":0,
"duration":102,
"bytes_in":1200,
"flow_id":"81454281200000731489",
"tenant":1}

我想基于一个共同的字段加入两者,例如“flow_id”,以防万一。 建议示例或方法。对 .fieldsGrouping 感到困惑,这是我用例的解决方案吗?

【问题讨论】:

    标签: apache-kafka apache-storm


    【解决方案1】:

    您可以使用 Tident API 进行连接:

    TridentTopology topology = new TridentTopology();
    // do some stuff here
    topology.join(stream1, new Fields("key"), stream2, new Fields("x"), new Fields("key", "a", "b", "c"));
    

    查看文档了解更多详情:https://storm.apache.org/releases/1.0.0/Trident-API-Overview.html

    如果你想使用低级API,使用fieldsGrouping是正确的(当然,你需要自己考虑“窗口化”)

    类似这样的:

    TopologyBuilder builder = new TopologyBuilder();
    builder.setSpout("spout1",...);
    builder.setSpout("spout2",...);
    
    builder.setSpout("join",...)
           .fieldsGrouping("spout1", new Fields("flow_id"))
           .fieldsGrouping("spout2", new Fields("flow_id"));
    

    【讨论】:

      猜你喜欢
      • 2014-09-02
      • 1970-01-01
      • 1970-01-01
      • 2018-08-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多