Flink TransForm


TransForm 简介

 转换算子可以把一个或多个DataStream转成一个新的DataStream.程序可以把多个复杂的转换组合成复杂的数据流拓扑。

 常用算子

1、map

作用 将数据流中的数据进行转换, 形成新的数据流,消费一个元素并产出一个元素

参数 lambda表达式或MapFunction实现类

返回 DataStream → DataStream

示例 

import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.wdh01.bean.WaterSensor;

import java.util.Arrays;
import java.util.List;

/**
 * 从集合读取数据
 */
public class Flink06_Transform_Map {
    public static void main(String[] args) throws Exception {
        //1、获取执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(2);//若 并行度1 按顺序读取数据,否则不按顺序
        // 2、从文件读取数据
        DataStreamSource stringDataStreamSource = env.readTextFile("input/sensor");
        //3、转换 javabean 打印数据
        stringDataStreamSource
                .map(new MyMapFunction())
                .print();

        //4、执行
        env.execute();

    }

    /**
     * 输入字符串,返回一个对象
     */
    public static class MyMapFunction implements MapFunction {
        @Override
        public WaterSensor map(String value) throws Exception {
            String[] split = value.split(",");
            return new WaterSensor(split[0], Long.parseLong(split[1]), Integer.parseInt(split[2]));
        }
    }
}
Rich...Function类 所有Flink函数类都有其Rich版本。它与常规函数的不同在于,可以获取运行环境的上下文,并拥有一些生命周期方法,所以可以实现更复杂的功能。也有意味着提供了更多的,更丰富的功能。例如:RichMapFunction
import org.apache.flink.api.common.functions.RichMapFunction;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.wdh01.bean.WaterSensor;

public class Flink07_Transform_RichMap {
    public static void main(String[] args) throws Exception {
        //1、获取执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);//若 并行度1 按顺序读取数据,否则不按顺序
        //2、读取数据
        DataStreamSource stringDataStreamSource = env.readTextFile("input/sensor");
        //3、将每行数据转换为 javabean
        SingleOutputStreamOperator map = stringDataStreamSource.map(new MyRichMapFunction());
        //4、打印
        map.print();
        //5、执行
        env.execute();
    }

    /**
     * RichMapFunction 有生命周期方法,可以获取上下文环境,作状态编程
     */
    public static class MyRichMapFunction extends RichMapFunction {

        @Override
        public WaterSensor map(String value) throws Exception {
            String[] split = value.split(",");
            return new WaterSensor(split[0], Long.parseLong(split[1]), Integer.parseInt(split[2]));
        }

        @Override //一个并行度 调用一次
        public void open(Configuration parameters) throws Exception {
            System.out.println(" open 执行一次 ");
        }

        @Override //一个并行度 调用两次
        public void close() throws Exception {
            System.out.println(" close 执行一次... ");
        }
    }

}

2、flatMap

作用 消费一个元素并产生零个或多个元素

参数  FlatMapFunction实现类

返回 DataStream → DataStream

示例

import org.apache.flink.api.common.functions.RichFlatMapFunction;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.util.Collector;

public class Flink08_Transform_RichFlatMap {
    public static void main(String[] args) throws Exception {
        //1、获取执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);//若 并行度1 按顺序读取数据,否则不按顺序
        //2、读取数据
        DataStreamSource stringDataStreamSource = env.readTextFile("input/sensor");
        //3、压平数据
        SingleOutputStreamOperator stringSingleOutputStreamOperator = stringDataStreamSource.flatMap(new MyRichFlatMap());
        //4、打印
        stringSingleOutputStreamOperator.print();
        //5、执行
        env.execute();
    }

    /**
     * 扁平化
     */
    public static class MyRichFlatMap extends RichFlatMapFunction {

        @Override
        public void flatMap(String value, Collector out) throws Exception {
            String[] split = value.split(",");
            for (String s : split) {
                out.collect(s);
            }
        }

        @Override
        public void close() throws Exception {
            super.close();
            System.out.println("close");
        }

        @Override
        public void open(Configuration parameters) throws Exception {
            super.open(parameters);
            System.out.println("open");
        }
    }
}

3、filter

作用 根据指定的规则将满足条件(true)的数据保留,不满足条件(false)的数据丢弃

参数 FlatMapFunction实现类

返回 DataStream → DataStream

示例
import org.apache.flink.api.common.functions.RichFilterFunction;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

public class Flink09_Transform_RichFilter {
    public static void main(String[] args) throws Exception {
        //1、获取执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);//若 并行度1 按顺序读取数据,否则不按顺序
        //2、读取数据
        DataStreamSource stringDataStreamSource = env.readTextFile("input/sensor");
        //3、过滤数据
        SingleOutputStreamOperator filter = stringDataStreamSource.filter(new MyRichFliterFunction());
        //4、打印
        filter.print();
        //5、执行
        env.execute();

    }

    /**
     * 过滤 过滤 水位小于30 的数据
     */
    public static class MyRichFliterFunction extends RichFilterFunction {

        @Override
        public boolean filter(String value) throws Exception {
            String[] split = value.split(",");
            //写出时进行过滤
            return Integer.parseInt(split[2]) > 30;
        }

        @Override
        public void close() throws Exception {
            super.close();
            System.out.println("close");
        }

        @Override
        public void open(Configuration parameters) throws Exception {
            super.open(parameters);
            System.out.println("open");
        }
    }
}

4、Connect

作用 在某些情况下,需要将两个不同来源的数据流进行连接,实现数据匹配,比如订单支付和第三方交易信息,这两个信息的数据就来自于不同数据源,连接后,将订单支付和第三方交易信息进行对账,此时,才能算真正的支付完成。Flink中的connect算子可以连接两个保持他们类型的数据流,两个数据流被connect之后,只是被放在了一个同一个流中,内部依然保持各自的数据和形式不发生任何变化,两个流相互独立。

 

参数   另外一个流

返回   DataStream[A], DataStream[B] -> ConnectedStreams[A,B]

import org.apache.flink.streaming.api.datastream.ConnectedStreams;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.co.CoMapFunction;

public class Flink10_Transform_Connect {
    public static void main(String[] args) throws Exception {
        //1、获取执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);//若 并行度1 按顺序读取数据,否则不按顺序
        //2、从端口读取数据
        DataStreamSource StringDS = env.socketTextStream("hadoop103", 9998);
        DataStreamSource socketTextStream2 = env.socketTextStream("hadoop103", 9988);
        //3、将 socketTextStream2 转换为 Int
        SingleOutputStreamOperator intDs = socketTextStream2.map(String::length); //方法引用
        //    SingleOutputStreamOperator map = socketTextStream2.map(data -> data.length()); //lambel 表达式
        //4、连接 2 个流
        ConnectedStreams connect = StringDS.connect(intDs);
       //5、处理连接之后的流
        SingleOutputStreamOperator map = connect.map(new CoMapFunction() {

            @Override
            public Object map1(String value) throws Exception {
                return value;
            }

            @Override
            public Object map2(Integer value) throws Exception {
                return value;
            }
        });
        //6、打印
        map.print();
        //7、执行
        env.execute();
    }
}


注意:
  1. 两个流中存储的数据类型可以不同
  2. 只是机械的合并在一起, 内部仍然是分离的2个流
  3. 只能2个流进行connect, 不能有第3个参与

5、union

 作用  对两个或者两个以上的DataStream进行union操作,产生一个包含所有DataStream元素的新DataStream

import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;


public class Flink01_Transform_Union {
    public static void main(String[] args) throws Exception {

        //1、获取执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        //2、读取端口数据
        DataStreamSource socketTextStream1 = env.socketTextStream("hadoop103", 9998);
        DataStreamSource socketTextStream2 = env.socketTextStream("hadoop103", 9988);
        //3、union 链接两个数据
        DataStream union = socketTextStream1.union(socketTextStream2);
        //4、打印
        union.print();
        //5、执行
        env.execute();
    }
}
connect与 union 区别:
  • union之前两个流的类型必须是一样,connect可以不一样
  • connect只能操作两个流,union可以操作多个。