【发布时间】:2018-03-02 18:07:05
【问题描述】:
我正在使用Flink 1.4.0。
假设我有一个POJO 如下:
public class Rating {
public String name;
public String labelA;
public String labelB;
public String labelC;
...
}
还有一个JOIN 函数:
public class SetLabelA implements JoinFunction<Tuple2<String, Rating>, Tuple2<String, String>, Tuple2<String, Rating>> {
@Override
public Tuple2<String, Rating> join(Tuple2<String, Rating> rating, Tuple2<String, String> labelA) {
rating.f1.setLabelA(labelA)
return rating;
}
}
假设我想应用JOIN 操作来设置DataSet<Tuple2<String, Rating>> 中每个字段的值,我可以这样做:
DataSet<Tuple2<String, Rating>> ratings = // [...]
DataSet<Tuple2<String, Double>> aLabels = // [...]
DataSet<Tuple2<String, Double>> bLabels = // [...]
DataSet<Tuple2<String, Double>> cLabels = // [...]
...
DataSet<Tuple2<String, Rating>>
newRatings =
ratings.leftOuterJoin(aLabels, JoinOperatorBase.JoinHint.REPARTITION_SORT_MERGE)
// key of the first input
.where("f0")
// key of the second input
.equalTo("f0")
// applying the JoinFunction on joining pairs
.with(new SetLabelA());
不幸的是,这是必要的,因为评级和所有xLabels 都非常大DataSets,我不得不查看每个xlabels 以找到我需要的字段值,同时它是并非每个xlabels 中都存在所有评级键。
这实际上意味着我必须为每个xlabel 执行leftOuterJoin,为此我还需要创建相应的JoinFunction 实现,该实现利用来自Rating POJO 的正确设置器。
有没有人能想到的更有效的方法来解决这个问题?
就分区策略而言,我确保将DataSet<Tuple2<String, Rating>> ratings 排序为:
DataSet<Tuple2<String, Rating>> sorted_ratings = ratings.sortPartition(0, Order.ASCENDING).setParallelism(1);
通过将并行度设置为 1,我可以确定整个数据集都会被排序。然后我使用.partitionByRange:
DataSet<Tuple2<String, Rating>> partitioned_ratings = sorted_ratings.partitionByRange(0).setParallelism(N);
N 是我的 VM 上的内核数。我在这里遇到的另一个问题是,第一个设置为 1 的.setParallelism 是否对管道其余部分的执行方式有限制,即后续.setParallelism(N) 是否可以改变DataSet 的处理方式?
最后,我做了所有这些,以便当partitioned_ratings 与xlabels DataSet 连接时,JOIN 操作将与JoinOperatorBase.JoinHint.REPARTITION_SORT_MERGE 一起完成。根据Flinkv.1.4.0 的文档:
REPARTITION_SORT_MERGE:系统对每个输入进行分区(打乱)(除非输入已经分区)并对每个输入进行排序(除非它已经排序)。输入通过排序输入的流合并连接。如果其中一个或两个输入都已排序,则此策略很好。
所以在我的情况下,ratings 已排序(我认为),而每个 xlabels DataSets 都没有,因此这是最有效的策略是有道理的。这有什么问题吗?任何替代方法?
【问题讨论】:
标签: java join apache-flink