【问题标题】:Does Storm Trident can only doing JOIN within batches?Storm Trident 只能分批做 JOIN 吗?
【发布时间】:2013-03-19 18:49:15
【问题描述】:

我想实现一个 JOIN 语义,我尝试了 Trident 拓扑中的 join 方法。 我发现加入是在批次之间进行的。 如果两个流之间的连接有数百万个元组,它必须在一个批次内吗?

在genderSpout中,每个batch有3个元组,所以Spout会发出2个batch ageSpout,每批次有 5 个元组,所以 Spout 只会发出 1 个批次

我用 JoinType 做一个 LEFT OUTER JOIN

测试代码的输出是:

1 man 15
2 woman 18
1 man 19
3 woman NULL
4 man NULL
1 woman NULL

从输出中,我发现前四个结果是连接来自genderSpout 的第一批和来自ageSpout 的第一批。 最后两个结果是来自 genderSpout 的第二批与来自 ageSpout 的空批之间的连接。 所以结果对于 JOIN 语义是不正确的,因为我想要的 genderSpout LEFT JOIN ageSpout 的结果是:

1 man 15
1 man 19
2 woman 18
3 woman NULL
4 man 20
1 woman 15
1 woman 19

所以我的问题是:如果JOIN的两边(Spout)有数百万个元组,我应该把它们放在一批中以获得正确的结果吗?

或者我走的路是错误的,你能告诉我我应该怎么做才能得到 OUTER JOIN 语义的正确结果?

测试代码如下:

public static void main(String[] args) throws Exception{
    Fields genderField = new Field("id", "gender");
    FixedBatchSpout genderSpout = new FixedBatchSpout(genderField, 3,
        new Values("1", "man"),
        new Values("2", "woman"),
        new Values("3", "woman"),
        new Values("4", "man"),
        new Values("1", "woman"));
    genderSpout.setCycle(false);

    Fields ageField = new Field("id2", "age");
    FiexedBatchSpout ageSpout = new FixedBatchSpout(new Fields("id2", "age"), 5,
        new Values("1", "15"),
        new Values("4", "20"),
        new Values("2", "18"),
        new Values("1", "19"));
    ageSpout.setCycle(false);

    List<Stream> allStreams = new ArrayList<Stream>();
    List<Fields> allFields = new ArrayList<Fields>();
    List<Fields> joinFileds = new ArrayList<Fields>();
    List<JoinType> joinTypes = new ArrayList<JoinType>();    

    TridentTopology topology = new TridentTopology();

    Stream genderStream = topology.newStream("genderIn", genderSpout);
    Stream ageStream = topology.newStream("ageIn", ageSpout);

    allStreams.add(genderStream);
    allStreams.add(ageStream);

    allFields.add(genderFields);
    allFields.add(ageFields);

    joinFields.add(new Field("id")));
    joinFields.add(new Field("id2"));

    joinTypes.add(JoinType.INNER);
    joinTypes.add(JoinType.OUTER);

    topology.join(allStreams, joinFields, new Filds("id", "gender", "age"), joinTypes)

    LocalCluster cluster = new LocalCluster();

    Config config = new Config()
    config.setDebug(false);
    config.setMaxSpoutPending(3);

    cluster.submitTopology("trident-join-test", config, topology.build());

    Thread.sleep(3000);
    cluster.shutdown();
}

【问题讨论】:

    标签: cluster-computing distributed apache-storm


    【解决方案1】:

    我在https://groups.google.com/forum/?fromgroups=#!forum/storm-user 中问过同样的问题: https://groups.google.com/forum/?fromgroups=#!topic/storm-user/7fxAVgF2_0M

    Jason Jackson 的回答是: TridentTopology.join 不会跨批次加入。您可以使用 stateQuery 和 partitionPersist 之一跨批次进行流式连接。

    希望有用

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2014-08-22
      • 2014-06-11
      • 2013-03-09
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-09-19
      • 1970-01-01
      相关资源
      最近更新 更多