【问题标题】:Partition Strategy for applying multiple JOINs on a Flink DataSet在 Flink DataSet 上应用多个 JOIN 的分区策略
【发布时间】: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&lt;Tuple2&lt;String, Rating&gt;&gt; 中每个字段的值,我可以这样做:

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&lt;Tuple2&lt;String, Rating&gt;&gt; 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_ratingsxlabels DataSet 连接时,JOIN 操作将与JoinOperatorBase.JoinHint.REPARTITION_SORT_MERGE 一起完成。根据Flinkv.1.4.0 的文档:

REPARTITION_SORT_MERGE:系统对每个输入进行分区(打乱)(除非输入已经分区)并对每个输入进行排序(除非它已经排序)。输入通过排序输入的流合并连接。如果其中一个或两个输入都已排序,则此策略很好。

所以在我的情况下,ratings 已排序(我认为),而每个 xlabels DataSets 都没有,因此这是最有效的策略是有道理的。这有什么问题吗?任何替代方法?

【问题讨论】:

    标签: java join apache-flink


    【解决方案1】:

    到目前为止,我还无法完成这个策略。似乎依赖JOINs 太麻烦了,因为它们是昂贵的操作,除非真的有必要,否则应该避免它们。

    例如,如果Datasets 的大小都非常大,则应使用JOINs。如果不是,一个方便的替代方法是使用BroadCastVariables,无论使用什么目的,两个Datasets(最小的)之一都会在工作人员之间广播。下面是一个示例(为方便起见,从link 复制)

    DataSet<Point> points = env.readCsv(...);
    
    DataSet<Centroid> centroids = ... ; // some computation
    
    points.map(new RichMapFunction<Point, Integer>() {
    
        private List<Centroid> centroids;
    
        @Override
        public void open(Configuration parameters) {
            this.centroids = getRuntimeContext().getBroadcastVariable("centroids");
        }
    
        @Override
        public Integer map(Point p) {
            return selectCentroid(centroids, p);
        }
    
    }).withBroadcastSet("centroids", centroids);
    

    此外,由于填充 POJO 的字段意味着将重复利用非常相似的代码,因此绝对应该使用jlens 以避免代码重复并编写更简洁且易于遵循的解决方案。

    【讨论】:

      猜你喜欢
      • 2022-08-21
      • 2015-03-25
      • 2020-04-30
      • 1970-01-01
      • 1970-01-01
      • 2019-06-22
      • 1970-01-01
      • 1970-01-01
      • 2013-07-14
      相关资源
      最近更新 更多