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、从文件读取数据 DataStreamSourceRich...Function类 所有Flink函数类都有其Rich版本。它与常规函数的不同在于,可以获取运行环境的上下文,并拥有一些生命周期方法,所以可以实现更复杂的功能。也有意味着提供了更多的,更丰富的功能。例如:RichMapFunctionstringDataStreamSource = 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])); } } }
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、读取数据 DataStreamSourcestringDataStreamSource = 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、读取数据 DataStreamSourcestringDataStreamSource = 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、读取数据 DataStreamSourcestringDataStreamSource = 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 个流 ConnectedStreamsconnect = StringDS.connect(intDs); //5、处理连接之后的流 SingleOutputStreamOperator
- 两个流中存储的数据类型可以不同
- 只是机械的合并在一起, 内部仍然是分离的2个流
- 只能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、读取端口数据 DataStreamSourceconnect与 union 区别:socketTextStream1 = env.socketTextStream("hadoop103", 9998); DataStreamSource socketTextStream2 = env.socketTextStream("hadoop103", 9988); //3、union 链接两个数据 DataStream union = socketTextStream1.union(socketTextStream2); //4、打印 union.print(); //5、执行 env.execute(); } }
- union之前两个流的类型必须是一样,connect可以不一样
- connect只能操作两个流,union可以操作多个。