我是靠谱客的博主 缥缈火龙果,最近开发中收集的这篇文章主要介绍Flink批处理-DataSet Transformations转化,觉得挺不错的,现在分享给大家,希望可以做个参考。

概述

环境

flink-1.9.0

 

一、需要的依赖

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-java</artifactId>
    <version>1.9.0</version>
</dependency>

二、初始化执行环境,读取用到的数据文件

/**
 * 获得执行环境,以字段的形式初始化,方便使用
 */
private static ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();

/**
 * 文本数据集
 */
private static DataSet<String> text = env.readTextFile("FlinkModule/src/main/resources/batch/text");

/**
 * 单词数据集
 */
private static DataSet<String> word = env.readTextFile("FlinkModule/src/main/resources/common/word");

文件text

This is a text file!
To be or not to be , this is a question!
end

文件word(每个单词之间按照制表符分隔)

how are    you
world  and    that
hello  world
jack   and
app    storm  storm  what
spark  spark

三、常用Transformation

说明:flink的方法可以通过匿名方法、Lambda表达式、类的形式实现转换函数,这里主要以匿名方法为主,简单使用了Lambda表达式和类的形式。对于以类的形式实现主要是为了提高复用性,只需要新建类实现需要实现的方法接口,重写方法即可。

1.map

使用匿名方法的方式实现

/**
 * map方法,匿名方法实现
 * 获取一个元素进行处理最后返回一个元素
 */
public static void map() throws Exception {
    text.map(new MapFunction<String, String>() {
        // 其中第一个泛型为输入类型,第二个泛型为返回值类型
        @Override
        public String map(String s) throws Exception {
            return "这一行的内容是:" + s;
        }
    // 调用print方法输出,如果只有数据集的转化,flink不会真正执行,
    // 只有触发需要shuffle的方法例如reduce才会真正执行此次任务。
    }).print();
}

使用Lambda表达式实现

/**
 * map方法,Lambda表达式实现
 * 获取一个元素进行处理最后返回一个元素
 */
public static void mapLambda() throws Exception {
    text.map((MapFunction<String, String>) s -> "这一行的内容是:" + s).print();
}

需要注意的是:使用Lambda表达式时有的返回值类型无法自动判断,需要显式说明,详细看flatMap的Lambda方式实现

运行结果:

这一行的内容是:To be or not to be , this is a question!
这一行的内容是:This is a text file!
这一行的内容是:end

2.flatMap

通过匿名方法实现

/**
 * flatMap方法,通过类去实现接口实现
 * 获取一个元素进行处理最后返回0个,一个或多个类型相同的元素
 */
public static void flatMap() throws Exception {
    word.flatMap(new FlatMapFunction<String, Tuple2<String,Integer>>() {
        // 将一行的单词切分成单个单词并形成元组收集起来
        @Override
        public void flatMap(String s, Collector<Tuple2<String, Integer>> collector) throws Exception {
            for(String str : s.split("t")){
                collector.collect(new Tuple2<>(str,1));
            }
        }
    }).print();
}

使用Lambda表达式实现

需要注意的是对于flatMap方法返回值类型为void,数据收集在collector收集器中,因此flink不能自动推断出数据的类型,需要显式的指明collector中数据的类型,即在flatMap后加上 returns(Types.TUPLE(Types.STRING,Types.INT))

/**
 * flatMap方法,Lambda表达式实现
 * 获取一个元素进行处理最后返回0个,一个或多个类型相同的元素
 */
public static void flatMapLambda() throws Exception {
    word.flatMap((FlatMapFunction<String, Tuple2<String, Integer>>) (s, collector) -> {
        for(String str : s.split("t")){
            collector.collect(new Tuple2<>(str,1));
        }
    }).returns(Types.TUPLE(Types.STRING,Types.INT)).print();
}

如果没有显示说明返回值类型则会报以下异常:

Exception in thread "main" org.apache.flink.api.common.functions.InvalidTypesException: The return type of function 'flatMapLamada(DataSetTransformations.java:75)' could not be determined automatically, due to type erasure. You can give type information hints by using the returns(...) method on the result of the transformation call, or by letting your function implement the 'ResultTypeQueryable' interface.
	at org.apache.flink.api.java.DataSet.getType(DataSet.java:178)
	at org.apache.flink.api.java.DataSet.collect(DataSet.java:410)
	at org.apache.flink.api.java.DataSet.print(DataSet.java:1652)
	at cn.myclass.batch.DataSetTransformations.flatMapLamada(DataSetTransformations.java:79)
	at cn.myclass.batch.DataSetTransformations.main(DataSetTransformations.java:86)
Caused by: org.apache.flink.api.common.functions.InvalidTypesException: The generic type parameters of 'Collector' are missing. In many cases lambda methods don't provide enough information for automatic type extraction when Java generics are involved. An easy workaround is to use an (anonymous) class instead that implements the 'org.apache.flink.api.common.functions.FlatMapFunction' interface. Otherwise the type has to be specified explicitly using type information.
	at org.apache.flink.api.java.typeutils.TypeExtractionUtils.validateLambdaType(TypeExtractionUtils.java:350)
	at org.apache.flink.api.java.typeutils.TypeExtractionUtils.extractTypeFromLambda(TypeExtractionUtils.java:176)
	at org.apache.flink.api.java.typeutils.TypeExtractor.getUnaryOperatorReturnType(TypeExtractor.java:571)
	at org.apache.flink.api.java.typeutils.TypeExtractor.getFlatMapReturnTypes(TypeExtractor.java:196)
	at org.apache.flink.api.java.DataSet.flatMap(DataSet.java:266)
	at cn.myclass.batch.DataSetTransformations.flatMapLamada(DataSetTransformations.java:75)
	... 1 more 

以类的形式实现flatMap

/**
 * 以类的形式实现flatMap方法,提高复用性
 * 将一行字符串按照制表符切分并形成元组
 */
public static class MyFlatMapFunction implements FlatMapFunction<String,Tuple2<String,Integer>> {

    @Override
    public void flatMap(String s, Collector<Tuple2<String,Integer>> collector) throws Exception {
        for(String str : s.split("t")){
            collector.collect(new Tuple2<>(str, 1));
        }
    }
}
/**
 * flatMap方法,匿名方法实现,获取一个元素进行处理最后返回0个,一个或多个类型相同的元素
 * 将一行单词按照制表符切分并形成元组放入收集器中
 */
public static void flatMapFunctionClass() throws Exception {
    word.flatMap(new MyFlatMapFunction()).print();
}

运行结果

(how,1)
(are,1)
(you,1)
(app,1)
(storm,1)
(storm,1)
(what,1)
(jack,1)
(and,1)
(spark,1)
(spark,1)
(hello,1)
(world,1)
(world,1)
(and,1)
(that,1)

3.mapPartition

分区进行map方法,数据会按照并行度分为n个区,然后各自执行此mapPartition方法

/**
 * mapPartition方法,分区进行map方法,分区个数取决于之前设置的并行度,默认为CPU核心数
 * 设置并行度并统计每个分区中的元素个数后输出,由于是并行计算所以输出时可能乱序
 */
public static void mapPartition() throws Exception {
    word.mapPartition(new MapPartitionFunction<String, Long>() {
        @Override
        public void mapPartition(Iterable<String> iterable, Collector<Long> collector) throws Exception {
            long c = 0;
            for (String s : iterable) {
                // 输出此行内容
                System.out.println(s);
                c++;

            }
            System.out.println("此分区有" + c + "个元素");
            collector.collect(c);
        }
    // 设置并行度为2,即有二个分区
    }).setParallelism(2).print();
}

运行结果(这里的并行度为2,由于是并行计算,所以输出结果为乱序,可以通过env.setParallelism(n)来设置并行度)

how	are	you
hello	world
world	and	that
jack	and
app	storm	storm	what
此分区有2个元素
spark	spark
此分区有4个元素
2
4

4.filter

过滤不符合条件的数据

/**
 * filter方法,按照条件过滤数据
 * 将字符串长度小于5的剔除
 */
public static void filter() throws Exception {
    text.filter(new FilterFunction<String>() {
        @Override
        public boolean filter(String s) throws Exception {
            return s.length()>5;
        }
    }).print();
}

运行结果

To be or not to be , this is a question!
This is a text file!
//因为end长度为3不符合条件,因此被过滤

5.project

投影元组中的字段,按照索引选取需要的字段

/**
 * project方法,投影元组中的字段,按照索引选取字段
 * 先将一行单词切分形成元组,再通过投影方法选取元组中的单词
 * Tuple2<String,Integer> --> Tuple<String>
 */
public static void project() throws Exception {
    DataSet tuple = word.flatMap(new MyFlatMapFunction());
    tuple.project(0).print();
}

运行结果

(world)
(and)
(that)
(spark)
(spark)
(how)
(are)
(you)
(hello)
(world)
(app)
(storm)
(storm)
(what)
(jack)
(and)

6.reduce

将元素进行归并,一般配合分组使用

/**
 * reduce方法,将元素进行归并,一般配合分组使用
 * 将字符串拼接到一起
 */
public static void reduce() throws Exception {
    text.reduce(new ReduceFunction<String>() {
        @Override
        public String reduce(String s1, String s2) throws Exception {
            // 将字符串拼接到一起
            return s1 + s2;
        }
    // 设置并发度为1
    }).setParallelism(1).print();
}

运行结果

endTo be or not to be , this is a question!This is a text file!

7.groupBy

将元素按照元组索引、Javabean字段、键选择器分组

/**
 * groupBy方法,将元素按照元组索引、Javabean字段、键选择器分组
 * 将单词切分后按照单词分组,统计每个单词出现的次数
 */
public static void groupBy() throws Exception {
    // 方法一,按照元组中的第一个索引列分组,即按照元组中的单词分组
    word.flatMap(new MyFlatMapFunction()).groupBy(0)
            .reduce(new ReduceFunction<Tuple2<String, Integer>>() {
                @Override
                public Tuple2<String, Integer> reduce(Tuple2<String, Integer> tuple1, Tuple2<String, Integer> tuple2) throws Exception {
                    // 将单词出现的次数累加
                    return new Tuple2<>(tuple1.f0, tuple1.f1 + tuple1.f1);
                }
            }).print();
    // 方法二、如果元素为javabean的话可以按照字段名分组,例如groupBy("word") 或者 groupBy("count")
    // 方法三、按照keySelector自定义分组,这里仍然以WordCount为例,以单词分组
    /*new KeySelector<Tuple2<String,Integer>, String>(){
        @Override
        public String getKey(Tuple2<String,Integer> tuple) throws Exception {
            return tuple.f0;
        }
    };*/
}

运行结果

(jack,1)
(what,1)
(you,1)
(world,2)
(hello,1)
(and,2)
(are,1)
(app,1)
(storm,2)
(how,1)
(spark,2)
(that,1)

8.groupReduce

对每个分组中的元组进行归并,可以返回一个或多个元素

/**
 * reduceGroup方法,对每个分组中的元组进行归并,可以返回一个或多个元素
 * 将相同分组的单词出现的次数进行累加
 */
public static void groupReduce() throws Exception {
    word.flatMap(new MyFlatMapFunction()).groupBy(0)
            .reduceGroup(new GroupReduceFunction<Tuple2<String, Integer>, Tuple2<String, Integer>>() {
                @Override
                public void reduce(Iterable<Tuple2<String, Integer>> iterable, Collector<Tuple2<String, Integer>> collector) throws Exception {
                    int count = 0;
                    String word = null;
                    for(Tuple2<String, Integer> tuple :iterable){
                        word = tuple.f0;
                        count += tuple.f1;
                    }
                    collector.collect(new Tuple2<>(word, count));
                }
            }).print();
}

运行结果

(jack,1)
(what,1)
(you,1)
(world,2)
(hello,1)
(and,2)
(are,1)
(app,1)
(storm,2)
(how,1)
(spark,2)
(that,1)

9.sortGroup

将分组按照某个字段或者元组的某个字段排序

/**
 * sortGroup方法,将分组按照某个字段或者元组的某个字段排序
 * 将单词元组按照出现次数分组后按照单词排序,最后取出前5个单词。
 * 这里由于单词初始次数为1,所以按照次数分组只有一个组,相当于对一个区内的单词排序,取出前5个元素
 */
public static void sortGroup() throws Exception {
    word.flatMap(new MyFlatMapFunction()).groupBy(1)
            .sortGroup(0, Order.ASCENDING)
            .reduceGroup(new GroupReduceFunction<Tuple2<String, Integer>, Tuple2<String, Integer>>() {
                @Override
                public void reduce(Iterable<Tuple2<String, Integer>> iterable, Collector<Tuple2<String, Integer>> collector) throws Exception {
                    Iterator<Tuple2<String, Integer>> iterator = iterable.iterator();
                    int size =0 ;
                    while (iterator.hasNext()) {
                        collector.collect(iterator.next());
                        if (++size == 5){
                            break;
                        }
                    }
                }
            }).setParallelism(1).print();
}

运行结果

(and,1)
(and,1)
(app,1)
(are,1)
(hello,1)

10.aggregate

将数据集按照某个字段进行聚合,可以求和、最小值、最大值

/**
 * aggregate方法,将数据集按照某个字段进行聚合,可以求和、最小值、最大值
 * 按照单词分组后统计单词出现次数的和,最小值,最大值
 */
public static void aggregate() throws Exception {
    word.flatMap(new MyFlatMapFunction()).groupBy(0)
            .aggregate(Aggregations.SUM,1).setParallelism(1).print();
    System.out.println("--------------");
    word.flatMap(new MyFlatMapFunction()).groupBy(0)
            .aggregate(Aggregations.MIN,1).setParallelism(1).print();
    System.out.println("--------------");
    word.flatMap(new MyFlatMapFunction()).groupBy(0)
            .aggregate(Aggregations.MAX,1).setParallelism(1).print();
}

运行结果

(and,2)
(app,1)
(are,1)
(hello,1)
(how,1)
(jack,1)
(spark,2)
(storm,2)
(that,1)
(what,1)
(world,2)
(you,1)
--------------
(and,1)
(app,1)
(are,1)
(hello,1)
(how,1)
(jack,1)
(spark,1)
(storm,1)
(that,1)
(what,1)
(world,1)
(you,1)
--------------
(and,1)
(app,1)
(are,1)
(hello,1)
(how,1)
(jack,1)
(spark,1)
(storm,1)
(that,1)
(what,1)
(world,1)
(you,1)

11.distinct

去除重复元素。默认去除数据集中的重复元素,可以根据元组(javabean)的字段、 表达式、选择器 选择符合条件的元素去除

/**
 * distinct方法,去除重复元素。默认去除数据集中的重复元素,可以根据元组(javabean)的字段、
 * 表达式、选择器 选择符合条件的元素去除
 * 将每一行的单词切分后去重
 */
public static void distinct() throws Exception {
    word.flatMap(new FlatMapFunction<String, String>() {
        @Override
        public void flatMap(String s, Collector<String> collector) throws Exception {
            for (String str:s.split("t")){
                collector.collect(str);
            }
        }
    }).distinct().print();
}

运行结果

jack
what
you
world
hello
and
are
app
storm
how
spark
that

12.join

将两个元组(javabean)进行连接,类似于数据库中的join连接。

/**
 * join方法,将两个元组(javabean)进行连接,类似于数据库中的join连接。
 * join后的where(0),equalTo(0)方法分别表示按照第一个元组的第一列和第二个元组的第一列连接,
 * 默认是按照两个字段的值等值连接,可以通过with()方法自定义连接方法,自定义方法分为两种,
 * 一种是JoinFunction,另一种是FlatJoinFunction,区别类似于map和flatMap,即前者是返回相同
 * 个元素,后者可以返回任何个元素
 *
 * 这里相当于以下两个文本按照单词进行连接,返回单词和单词出现的总数
 * 注意:左侧的每个元素都会与右侧符合条件的元素进行连接
 * (how,1)    连接      (and,2)
 * (are,1)              (app,1)
 * (you,1)              (are,1)
 * (app,1)              (hello,1)
 * (storm,1)            (how,1)
 * (storm,1)            (jack,1)
 * (what,1)             (spark,2)
 * (jack,1)             (storm,2)
 * (and,1)              (that,1)
 * (spark,1)            (what,1)
 * (spark,1)            (world,2)
 * (hello,1)            (you,1)
 * (world,1)
 * (world,1)
 * (and,1)
 * (that,1)
 */
public static void join() throws Exception {
    word.flatMap(new MyFlatMapFunction())
            .join(word.flatMap(new MyFlatMapFunction()).groupBy(0).aggregate(Aggregations.SUM,1))
            .where(0)
            .equalTo(0)
            .with(new JoinFunction<Tuple2<String, Integer>, Tuple2<String, Integer>, Tuple2<String, Integer>>() {
                @Override
                public Tuple2<String, Integer> join(Tuple2<String, Integer> tuple1, Tuple2<String, Integer> tuple2) throws Exception {

                    return new Tuple2<>(tuple1.f0, tuple1.f1+tuple2.f1);
                }
            }).print();
}

运行结果

(jack,2)
(what,2)
(you,2)
(world,3)
(world,3)
(hello,2)
(and,3)
(and,3)
(are,2)
(app,2)
(storm,3)
(storm,3)
(how,2)
(spark,3)
(spark,3)
(that,2)

13.OuterJoin

外连接,分为左外连接,右外连接,全外连接。与内连接不同的是外连接会保留主表的所有数据以及连接表中符合条件的数据

/**
 * OuterJoin方法,外连接,类似于数据库中的外连接,分为左外连接,右外连接,全外连接。
 * 与内连接不同的是外连接会保留主表的所有数据以及连接表中符合条件的数据。
 * 
 * 这里是右外连接
 * (how,1)   右外连接   (and,2)
 * (are,1)              (app,1)
 * (you,1)              (are,1)
 * (app,1)              (hello,1)
 * (storm,1)            (how,1)
 * (storm,1)            (jack,1)
 * (what,1)             (spark,2)
 * (jack,1)             (storm,2)
 * (and,1)              (that,1)
 * (spark,1)            (what,1)
 * (spark,1)            (world,2)
 * (hello,1)            (you,1)
 * (world,1)
 * (world,1)
 * (and,1)
 * (that,1)
 */
public static void outerJoin() throws Exception {
    word.flatMap(new MyFlatMapFunction())
            .rightOuterJoin(word.flatMap(new MyFlatMapFunction()).groupBy(0).aggregate(Aggregations.SUM,1))
            .where(0)
            .equalTo(0)
            .with(new JoinFunction<Tuple2<String, Integer>, Tuple2<String, Integer>, Tuple2<String, Integer>>() {
                @Override
                public Tuple2<String, Integer> join(Tuple2<String, Integer> tuple1, Tuple2<String, Integer> tuple2) throws Exception {

                    return new Tuple2<>(tuple1.f0, tuple1.f1+tuple2.f1);
                }
            }).print();
}

运行结果

(jack,2)
(what,2)
(you,2)
(world,3)
(world,3)
(hello,2)
(and,3)
(and,3)
(are,2)
(app,2)
(storm,3)
(storm,3)
(how,2)
(spark,3)
(spark,3)
(that,2)

14.Cross

求两个数据集的笛卡儿积

/**
 * cross方法,求两个数据集的笛卡儿积,可以配合with方法使用,自定义连接方法
 * 将text文本与自身进行笛卡儿积
 */
public static void cross() throws Exception {
    text.cross(text).print();
}

运行结果

(This is a text file!,To be or not to be , this is a question!)
(This is a text file!,This is a text file!)
(This is a text file!,end)
(end,To be or not to be , this is a question!)
(end,This is a text file!)
(end,end)
(To be or not to be , this is a question!,To be or not to be , this is a question!)
(To be or not to be , this is a question!,This is a text file!)
(To be or not to be , this is a question!,end)

15.union

将连个数据集连接到一起形成新的数据集

/**
 * union方法,将两个数据集连接到一起形成新的数据集
 */
public static void union() throws Exception {
    text.union(text).print();
}

运行结果

To be or not to be , this is a question!
To be or not to be , this is a question!
This is a text file!
This is a text file!
end
end

16.coGroup

协分组,将两个数据集分组后将相同分组进行连接执行一系列操作

/**
 * coGroup方法,协分组,将两个数据集分组后将相同分组进行连接执行一系列操作
 * 将单词数据集与自己协分组然后统计出现的次数
 */
public static void coGroup() throws Exception {
    word.flatMap(new MyFlatMapFunction())
            .coGroup(word.flatMap(new MyFlatMapFunction()))
            .where(0)
            .equalTo(0)
            .with(new CoGroupFunction<Tuple2<String, Integer>, Tuple2<String, Integer>, Tuple2<String, Integer>>() {
                @Override
                public void coGroup(Iterable<Tuple2<String, Integer>> iterable, Iterable<Tuple2<String, Integer>> iterable1, Collector<Tuple2<String, Integer>> collector) throws Exception {
                    String word = null;
                    int count = 0;
                    for(Tuple2<String,Integer> t:iterable){
                        word = t.f0;
                        count += t.f1;
                    }
                    for(Tuple2<String,Integer> t:iterable1){
                        count += t.f1;
                    }
                    collector.collect(new Tuple2<>(word,count));
                }
            }).print();
}

运行结果

(jack,2)
(what,2)
(you,2)
(world,4)
(hello,2)
(and,4)
(are,2)
(app,2)
(storm,4)
(how,2)
(spark,4)
(that,2)

17.first

从数据集中获得前n条数据,可以用于常规、分组、分区数据集

/**
 * first方法,从数据集中获得前n条数据,可以用于常规、分组、分区数据集
 * 输出text文本中的第一行数据
 */
public static void first() throws Exception {
    text.first(1).print();
}

运行结果

To be or not to be , this is a question!

以下为分区策略,分区后简单输出结果

18.rebalance

将数据集均匀的重新分配到每一个分区

/**
 * rebalance方法,将数据集均匀的重新分配到每一个分区
 */
public static void rebalance() throws Exception {
    text.rebalance().setParallelism(2).print();
}

19.hashPartition

将数据集按照某个字段的hash值分区

/**
 * partitionByHash方法,将数据集按照某个字段的hash值分区
 */
public static void hashPartition() throws Exception {
    word.flatMap(new MyFlatMapFunction()).partitionByHash(0).print();
}

20.rangePartition

将数据集按照一定范围分区

/**
 * partitionByRange方法,将数据集按照一定范围分区
 */
public static void rangePartition() throws Exception {
    word.flatMap(new MyFlatMapFunction()).partitionByRange(0).print();
}

21.sortPartition

将数据集按照某个字段分区并排序

/**
 * sortPartition方法,将数据集按照某个字段分区并排序
 */
public static void sortPartition() throws Exception {
    word.flatMap(new MyFlatMapFunction())
            .sortPartition(0, Order.ASCENDING)
            .sortPartition(1, Order.DESCENDING).print();
}

四、完整代码

package cn.myclass.batch;

import org.apache.flink.api.common.functions.*;
import org.apache.flink.api.common.operators.Order;
import org.apache.flink.api.common.typeinfo.Types;
import org.apache.flink.api.java.DataSet;
import org.apache.flink.api.java.ExecutionEnvironment;
import org.apache.flink.api.java.aggregation.Aggregations;
import org.apache.flink.api.java.functions.KeySelector;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.util.Collector;

import java.util.Iterator;


/**
 * 数据集变换
 * @author Yang
 */
public class DataSetTransformations {

    /**
     * 获得执行环境
     */
    private static ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();

    /**
     * 文本数据集
     */
    private static DataSet<String> text = env.readTextFile("FlinkModule/src/main/resources/batch/text");

    /**
     * 单词数据集
     */
    private static DataSet<String> word = env.readTextFile("FlinkModule/src/main/resources/common/word");

    /**
     * 以类的形式实现flatMap方法,提高复用性
     * 将一行字符串按照制表符切分并形成元组
     */
    public static class MyFlatMapFunction implements FlatMapFunction<String,Tuple2<String,Integer>> {

        @Override
        public void flatMap(String s, Collector<Tuple2<String,Integer>> collector) throws Exception {
            for(String str : s.split("t")){
                collector.collect(new Tuple2<>(str, 1));
            }
        }
    }

    /**
     * map方法,匿名方法实现,获取一个元素进行处理最后返回一个元素
     * 输出此行文本
     */
    public static void map() throws Exception {
        // 其中第一个泛型为输入类型,第二个泛型为返回值类型
        text.map(new MapFunction<String, String>() {
            @Override
            public String map(String s) throws Exception {
                return "这一行的内容是:" + s;
            }
        // 调用print方法输出。如果只有数据集的转化,flink不会真正执行,
        // 只有触发需要shuffle的方法,例如reduce才会真正执行此次任务。
        }).print();
    }

    /**
     * map方法,Lambda表达式实现
     * 输出此行文本
     */
    public static void mapLambda() throws Exception {
        text.map((MapFunction<String, String>) s -> "这一行的内容是:" + s).print();
    }

    /**
     * flatMap方法,匿名方法实现,获取一个元素进行处理最后返回0个,一个或多个类型相同的元素
     * 将一行单词按照制表符切分并形成元组放入收集器中
     */
    public static void flatMap() throws Exception {
        word.flatMap(new FlatMapFunction<String, Tuple2<String,Integer>>() {
            @Override
            public void flatMap(String s, Collector<Tuple2<String, Integer>> collector) throws Exception {
                // 将一行的单词切分成单个单词并形成元组返回
                for(String str : s.split("t")){
                    collector.collect(new Tuple2<>(str, 1));
                }
            }
        }).print();
    }

    /**
     * flatMap方法,Lambda表达式实现
     * 将一行单词按照制表符切分并形成元组放入收集器中
     */
    public static void flatMapLambda() throws Exception {
        word.flatMap((FlatMapFunction<String, Tuple2<String, Integer>>) (s, collector) -> {
            for(String str : s.split("t")){
                collector.collect(new Tuple2<>(str,1));
            }
        }).returns(Types.TUPLE(Types.STRING,Types.INT)).print();
    }

    /**
     * flatMap方法,匿名方法实现,获取一个元素进行处理最后返回0个,一个或多个类型相同的元素
     * 将一行单词按照制表符切分并形成元组放入收集器中
     */
    public static void flatMapFunctionClass() throws Exception {
        word.flatMap(new MyFlatMapFunction()).print();
    }


    /**
     * mapPartition方法,分区进行map方法,分区个数取决于之前设置的并行度,默认为CPU核心数
     * 设置并行度并统计每个分区中的元素个数后输出。由于是并行计算所以输出时可能乱序
     */
    public static void mapPartition() throws Exception {
        word.mapPartition(new MapPartitionFunction<String, Long>() {
            @Override
            public void mapPartition(Iterable<String> iterable, Collector<Long> collector) throws Exception {
                long c = 0;
                for (String s : iterable) {
                    // 输出此行内容
                    System.out.println(s);
                    c++;

                }
                System.out.println("此分区有" + c + "个元素");
                collector.collect(c);
            }
        // 设置并行度为2,即有二个分区
        }).setParallelism(2).print();
    }

    /**
     * filter方法,按照条件过滤数据
     * 将字符串长度小于5的数据剔除
     */
    public static void filter() throws Exception {
        text.filter(new FilterFunction<String>() {
            @Override
            public boolean filter(String s) throws Exception {
                return s.length()>5;
            }
        }).print();
    }
    
    /**
     * project方法,投影元组中的字段,按照索引选取字段
     * 先将一行单词切分形成元组,再通过投影方法选取元组中的单词
     * Tuple2<String,Integer> --> Tuple<String>
     */
    public static void project() throws Exception {
        DataSet tuple = word.flatMap(new MyFlatMapFunction());
        tuple.project(0).print();
    }

    /**
     * reduce方法,将元素进行归并,一般配合分组使用
     * 将字符串拼接到一起
     */
    public static void reduce() throws Exception {
        text.reduce(new ReduceFunction<String>() {
            @Override
            public String reduce(String s1, String s2) throws Exception {
                // 将字符串拼接到一起
                return s1 + s2;
            }
        // 设置并发度为1
        }).setParallelism(1).print();
    }

    /**
     * groupBy方法,按照元组索引、Javabean字段、选择器分组
     * 将单词切分后按照单词分组,统计每个单词出现的次数
     */
    public static void groupBy() throws Exception {
        // 方法一,按照元组中的第一个索引列分组,即按照元组中的单词分组
        word.flatMap(new MyFlatMapFunction()).groupBy(0)
                .reduce(new ReduceFunction<Tuple2<String, Integer>>() {
                    @Override
                    public Tuple2<String, Integer> reduce(Tuple2<String, Integer> tuple1, Tuple2<String, Integer> tuple2) throws Exception {
                        // 将单词出现的次数累加
                        return new Tuple2<>(tuple1.f0, tuple1.f1 + tuple1.f1);
                    }
                }).print();
        // 方法二、如果元素为javabean的话可以按照字段名分组,例如groupBy("word") 或者 groupBy("count")
        // 方法三、按照keySelector自定义分组,这里仍然以WordCount为例,以单词分组
        /*new KeySelector<Tuple2<String,Integer>, String>(){
            @Override
            public String getKey(Tuple2<String,Integer> tuple) throws Exception {
                return tuple.f0;
            }
        };*/
    }

    /**
     * reduceGroup方法,对每个分组中的元组进行归并,可以返回一个或多个元素
     * 将相同分组的单词出现的次数进行累加
     */
    public static void groupReduce() throws Exception {
        word.flatMap(new MyFlatMapFunction()).groupBy(0)
                .reduceGroup(new GroupReduceFunction<Tuple2<String, Integer>, Tuple2<String, Integer>>() {
                    @Override
                    public void reduce(Iterable<Tuple2<String, Integer>> iterable, Collector<Tuple2<String, Integer>> collector) throws Exception {
                        int count = 0;
                        String word = null;
                        for(Tuple2<String, Integer> tuple :iterable){
                            word = tuple.f0;
                            count += tuple.f1;
                        }
                        collector.collect(new Tuple2<>(word, count));
                    }
                }).print();
    }

    /**
     * sortGroup方法,将分组按照某个字段或者元组的某个字段排序
     * 将单词元组按照出现次数分组后按照单词排序,最后取出前5个单词。
     * 这里由于单词初始次数为1,所以按照次数分组只有一个组,相当于对一个区内的单词排序,取出前5个元素
     */
    public static void sortGroup() throws Exception {
        word.flatMap(new MyFlatMapFunction()).groupBy(1)
                .sortGroup(0, Order.ASCENDING)
                // 每个组调用一次reduce方法
                .reduceGroup(new GroupReduceFunction<Tuple2<String, Integer>, Tuple2<String, Integer>>() {
                    @Override
                    public void reduce(Iterable<Tuple2<String, Integer>> iterable, Collector<Tuple2<String, Integer>> collector) throws Exception {
                        Iterator<Tuple2<String, Integer>> iterator = iterable.iterator();
                        int size =0 ;
                        // 迭代器中的元素为一个组中的元素
                        while (iterator.hasNext()) {
                            collector.collect(iterator.next());
                            if (++size == 5){
                                break;
                            }
                        }
                    }
                }).setParallelism(1).print();
    }

    /**
     * aggregate方法,将数据集按照某个字段进行聚合,可以求和、最小值、最大值
     * 按照单词分组后统计单词出现次数的和,最小值,最大值
     */
    public static void aggregate() throws Exception {
        word.flatMap(new MyFlatMapFunction()).groupBy(0)
                .aggregate(Aggregations.SUM,1).setParallelism(1).print();
        System.out.println("--------------");
        word.flatMap(new MyFlatMapFunction()).groupBy(0)
                .aggregate(Aggregations.MIN,1).setParallelism(1).print();
        System.out.println("--------------");
        word.flatMap(new MyFlatMapFunction()).groupBy(0)
                .aggregate(Aggregations.MAX,1).setParallelism(1).print();
    }

    /**
     * distinct方法,去除重复元素。默认去除数据集中的重复元素,可以根据元组(javabean)的字段、
     * 表达式、选择器 选择符合条件的元素去除
     * 将每一行的单词切分后去重
     */
    public static void distinct() throws Exception {
        word.flatMap(new FlatMapFunction<String, String>() {
            @Override
            public void flatMap(String s, Collector<String> collector) throws Exception {
                for (String str:s.split("t")){
                    collector.collect(str);
                }
            }
        }).distinct().print();
    }

    /**
     * join方法,将两个元组(javabean)进行连接,类似于数据库中的join连接。
     * join后的where(0),equalTo(0)方法分别表示按照第一个元组的第一列和第二个元组的第一列连接,
     * 默认是按照两个字段的值等值连接,可以通过with()方法自定义连接方法,自定义方法分为两种,
     * 一种是JoinFunction,另一种是FlatJoinFunction,区别类似于map和flatMap,即前者是返回相同
     * 个元素,后者可以返回任何个元素
     *
     * 这里相当于以下两个文本按照单词进行连接,返回单词和单词出现的总数。
     * 注意:左侧的每个元素都会与右侧符合条件的元素进行连接
     * (how,1)    连接      (and,2)
     * (are,1)              (app,1)
     * (you,1)              (are,1)
     * (app,1)              (hello,1)
     * (storm,1)            (how,1)
     * (storm,1)            (jack,1)
     * (what,1)             (spark,2)
     * (jack,1)             (storm,2)
     * (and,1)              (that,1)
     * (spark,1)            (what,1)
     * (spark,1)            (world,2)
     * (hello,1)            (you,1)
     * (world,1)
     * (world,1)
     * (and,1)
     * (that,1)
     */
    public static void join() throws Exception {
        word.flatMap(new MyFlatMapFunction())
                .join(word.flatMap(new MyFlatMapFunction()).groupBy(0).aggregate(Aggregations.SUM,1))
                .where(0)
                .equalTo(0)
                .with(new JoinFunction<Tuple2<String, Integer>, Tuple2<String, Integer>, Tuple2<String, Integer>>() {
                    @Override
                    public Tuple2<String, Integer> join(Tuple2<String, Integer> tuple1, Tuple2<String, Integer> tuple2) throws Exception {

                        return new Tuple2<>(tuple1.f0, tuple1.f1+tuple2.f1);
                    }
                }).print();
    }

    /**
     * OuterJoin方法,外连接,类似于数据库中的外连接,分为左外连接,右外连接,全外连接。
     * 与内连接不同的是外连接会保留主表的所有数据以及连接表中符合条件的数据。
     *
     * 这里是右外连接
     * (how,1)   右外连接   (and,2)
     * (are,1)              (app,1)
     * (you,1)              (are,1)
     * (app,1)              (hello,1)
     * (storm,1)            (how,1)
     * (storm,1)            (jack,1)
     * (what,1)             (spark,2)
     * (jack,1)             (storm,2)
     * (and,1)              (that,1)
     * (spark,1)            (what,1)
     * (spark,1)            (world,2)
     * (hello,1)            (you,1)
     * (world,1)
     * (world,1)
     * (and,1)
     * (that,1)
     */
    public static void outerJoin() throws Exception {
        word.flatMap(new MyFlatMapFunction())
                .rightOuterJoin(word.flatMap(new MyFlatMapFunction()).groupBy(0).aggregate(Aggregations.SUM,1))
                .where(0)
                .equalTo(0)
                .with(new JoinFunction<Tuple2<String, Integer>, Tuple2<String, Integer>, Tuple2<String, Integer>>() {
                    @Override
                    public Tuple2<String, Integer> join(Tuple2<String, Integer> tuple1, Tuple2<String, Integer> tuple2) throws Exception {

                        return new Tuple2<>(tuple1.f0, tuple1.f1+tuple2.f1);
                    }
                }).print();
    }

    /**
     * cross方法,求两个数据集的笛卡儿积,可以配合with方法使用,自定义连接方法
     * 将text文本与自身进行笛卡儿积
     */
    public static void cross() throws Exception {
        text.cross(text).print();
    }

    /**
     * union方法,将两个数据集连接到一起形成新的数据集
     */
    public static void union() throws Exception {
        text.union(text).print();
    }

    /**
     * coGroup方法,协分组,将两个数据集分组后将相同分组进行连接执行一系列操作
     * 将单词数据集与自己协分组然后统计出现的次数
     */
    public static void coGroup() throws Exception {
        word.flatMap(new MyFlatMapFunction())
                .coGroup(word.flatMap(new MyFlatMapFunction()))
                .where(0)
                .equalTo(0)
                .with(new CoGroupFunction<Tuple2<String, Integer>, Tuple2<String, Integer>, Tuple2<String, Integer>>() {
                    @Override
                    public void coGroup(Iterable<Tuple2<String, Integer>> iterable, Iterable<Tuple2<String, Integer>> iterable1, Collector<Tuple2<String, Integer>> collector) throws Exception {
                        String word = null;
                        int count = 0;
                        for(Tuple2<String,Integer> t:iterable){
                            word = t.f0;
                            count += t.f1;
                        }
                        for(Tuple2<String,Integer> t:iterable1){
                            count += t.f1;
                        }
                        collector.collect(new Tuple2<>(word,count));
                    }
                }).print();
    }

    /**
     * first方法,从数据集中获得前n条数据,可以用于常规、分组、分区数据集
     * 输出text文本中的第一行数据
     */
    public static void first() throws Exception {
        text.first(1).print();
    }

    /**
     * rebalance方法,将数据集均匀的重新分配到每一个分区
     */
    public static void rebalance() throws Exception {
        text.rebalance().setParallelism(2).print();
    }

    /**
     * partitionByHash方法,将数据集按照某个字段的hash值分区
     */
    public static void hashPartition() throws Exception {
        word.flatMap(new MyFlatMapFunction()).partitionByHash(0).print();
    }

    /**
     * partitionByRange方法,将数据集按照一定范围分区
     */
    public static void rangePartition() throws Exception {
        word.flatMap(new MyFlatMapFunction()).partitionByRange(0).print();
    }

    /**
     * sortPartition方法,将数据集按照某个字段分区并排序
     */
    public static void sortPartition() throws Exception {
        word.flatMap(new MyFlatMapFunction())
                .sortPartition(0, Order.ASCENDING)
                .sortPartition(1, Order.DESCENDING).print();
    }


    public static void main(String[] args) throws Exception {
        map();
        mapLambda();
        flatMap();
        flatMapLambda();
        flatMapFunctionClass();
        mapPartition();
        filter();
        project();
        reduce();
        groupBy();
        groupReduce();
        sortGroup();
        aggregate();
        distinct();
        join();
        outerJoin();
        cross();
        union();
        coGroup();
        rebalance();
        hashPartition();
        rangePartition();
        sortPartition();
        first();
    }
}

如有错误,望指正!

 

最后

以上就是缥缈火龙果为你收集整理的Flink批处理-DataSet Transformations转化的全部内容,希望文章能够帮你解决Flink批处理-DataSet Transformations转化所遇到的程序开发问题。

如果觉得靠谱客网站的内容还不错,欢迎将靠谱客网站推荐给程序员好友。

本图文内容来源于网友提供,作为学习参考使用,或来自网络收集整理,版权属于原作者所有。
点赞(58)

评论列表共有 0 条评论

立即
投稿
返回
顶部