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、读取数据
        DataSource input = 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、读取数据
        DataStreamSource input = 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、读取端口数据
        DataStreamSource socketTextStream = 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 服务,在运行应用程序,否则运行程序直接报错