Flink 实现 WordCount
pom.xml
<properties> <flink.version>1.12.0flink.version> <java.version>1.8java.version> <scala.binary.version>2.11scala.binary.version> <slf4j.version>1.7.30slf4j.version> properties> <dependencies> <dependency> <groupId>org.apache.flinkgroupId> <artifactId>flink-javaartifactId> <version>${flink.version}version> dependency> <dependency> <groupId>org.apache.flinkgroupId> <artifactId>flink-streaming-java_${scala.binary.version}artifactId> <version>${flink.version}version> dependency> <dependency> <groupId>org.apache.flinkgroupId> <artifactId>flink-clients_${scala.binary.version}artifactId> <version>${flink.version}version> dependency> <dependency> <groupId>org.apache.flinkgroupId> <artifactId>flink-runtime-web_${scala.binary.version}artifactId> <version>${flink.version}version> dependency> <dependency> <groupId>org.slf4jgroupId> <artifactId>slf4j-apiartifactId> <version>${slf4j.version}version> dependency> <dependency> <groupId>org.slf4jgroupId> <artifactId>slf4j-log4j12artifactId> <version>${slf4j.version}version> dependency> <dependency> <groupId>org.apache.logging.log4jgroupId> <artifactId>log4j-to-slf4jartifactId> <version>2.14.0version> dependency> dependencies>
flink 批处理 WordCount
import org.apache.flink.api.common.functions.FlatMapFunction; import org.apache.flink.api.common.functions.MapFunction; import org.apache.flink.api.java.ExecutionEnvironment; import org.apache.flink.api.java.operators.*; import org.apache.flink.api.java.tuple.Tuple2; import org.apache.flink.util.Collector; public class Flink01_WordCount_Batch_0210 { public static void main(String[] args) throws Exception { //1、获取执行环境 ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment(); //2、读取数据 DataSourceinput = env.readTextFile("input"); //3、压平 FlatMapOperator wordDS = input.flatMap(new MyFlatMapFunction()); //4、将单词转换成元祖 MapOperator > wordToOnwDS = wordDS.map(new MapFunction >() { public Tuple2 map(String value) throws Exception { return new Tuple2 (value, 1); } }); //5、分组 按 MapOperator > 的 key 进行分组 UnsortedGrouping> groupByDS = wordToOnwDS.groupBy(0); //6、聚合 UnsortedGrouping > 聚合第二个字段 AggregateOperator> sum = groupByDS.sum(1); //7、打印 sum.print(); } /** * 自定义实现 flatMap 的 入参函数实现 压平操作 */ public static class MyFlatMapFunction implements FlatMapFunction { public void flatMap(String value, Collector out) throws Exception { //按照空格切分数据 String[] words = value.split(" "); //写出单词 for (String word : words) { out.collect(word); } } } }
准备好样例数据,运行程序可以看到运行效果。
flink 有界流 woedcount
import org.apache.flink.api.common.functions.FlatMapFunction; import org.apache.flink.api.java.functions.KeySelector; import org.apache.flink.api.java.tuple.Tuple2; import org.apache.flink.streaming.api.datastream.DataStreamSource; import org.apache.flink.streaming.api.datastream.KeyedStream; import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.util.Collector; /** * 有界流 wc */ public class Flink02_WordCount_Bounded_0210 { public static void main(String[] args) throws Exception { //1、获取执行环境 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); //设置并行度 env.setParallelism(1); //2、读取数据 DataStreamSourceinput = env.readTextFile("input"); //3、压平 SingleOutputStreamOperator > wordToone = input.flatMap(new LineToTupleFlatMapFunction()); //4、分组 KeyedStream , Object> keyStream = wordToone.keyBy(new KeySelector , Object>() { public Object getKey(Tuple2 value) throws Exception { return value.f0; } }); //5、分组 SingleOutputStreamOperator > sum = keyStream.sum(1); //6、打印 sum.print(); //7、启动任务 env.execute(); } /** * 压平&转换元组 */ public static class LineToTupleFlatMapFunction implements FlatMapFunction > { public void flatMap(String value, Collector > out) throws Exception { //切分数据 String[] words = value.split(" "); //遍历写出 for (String word : words) { out.collect(new Tuple2 (word, 1)); } } } }
flink 无界流 woedcount
import org.apache.flink.api.java.functions.KeySelector; import org.apache.flink.api.java.tuple.Tuple2; import org.apache.flink.streaming.api.datastream.DataStreamSource; import org.apache.flink.streaming.api.datastream.KeyedStream; import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; /** * 无界流 */ public class Flink_03_WordCount_UnBounded_0210 { /** * 并行度优先级 * 1、代码中算子单独设置 sum.print("sum --> ").setParallelism(3); * 2、代码中 env 全局设置 * 3、提交参数 * 4、默认配置 * * @param args * @throws Exception */ public static void main(String[] args) throws Exception { //1、获取执行环境 StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); //设置并行度 env.setParallelism(4); //2、读取端口数据 DataStreamSourcesocketTextStream = env.socketTextStream("hadoop103", 9998); //3、压平 转换 元组 SingleOutputStreamOperator > wordtoOne = socketTextStream.flatMap(new Flink02_WordCount_Bounded_0210.LineToTupleFlatMapFunction()).setParallelism(2); //4、分组 KeyedStream , Object> keyedStream = wordtoOne.keyBy(new KeySelector , Object>() { public Object getKey(Tuple2 value) throws Exception { return value.f0; } }); //5、聚合 SingleOutputStreamOperator > sum = keyedStream.sum(1); //6、打印测试 /* socketTextStream.print("line --> "); wordtoOne.print("wordtoOne --> ");*/ sum.print("sum --> ").setParallelism(3); env.execute(); //测试 hadoop103 执行 输入 nc -lk 9998 进入 阻塞状态 输入单词即可 } }
注意:需要提前运行 nc 服务,在运行应用程序,否则运行程序直接报错