依赖配置

 <properties>
    <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
    <maven.compiler.source>1.8</maven.compiler.source>
    <maven.compiler.target>1.8</maven.compiler.target>
    <flink.version>1.16.0</flink.version>
    <slf4j.version>1.7.36</slf4j.version>
    <log4j.version>2.17.2</log4j.version>
  </properties>

  <dependencies>
    <!-- Flink批和流开发依赖包 -->
    <dependency>
      <groupId>org.apache.flink</groupId>
      <artifactId>flink-clients</artifactId>
      <version>${flink.version}</version>
    </dependency>

    <!-- slf4j&log4j 日志相关包 -->
    <dependency>
      <groupId>org.slf4j</groupId>
      <artifactId>slf4j-log4j12</artifactId>
      <version>${slf4j.version}</version>
    </dependency>
    <dependency>
      <groupId>org.apache.logging.log4j</groupId>
      <artifactId>log4j-to-slf4j</artifactId>
      <version>${log4j.version}</version>
    </dependency>

  </dependencies>

单词文件

在这里插入图片描述

批量处理

public class BatchWordCount {
    public static void main(String[] args) throws Exception {
        // 1. 准备环境
        ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();

        // 2. 从文件读取数据
        DataSource<String> lineDS = env.readTextFile("./data/words.txt");

        // 3. 对数据进行切分操作
        FlatMapOperator<String, String> wordsDS = lineDS
                .flatMap((FlatMapFunction<String, String>) (line, out) -> {
            for (String word : line.split(" ")) {
                out.collect(word);
            }
        }).returns(Types.STRING);

        // 5. 对单词进行计数
        MapOperator<String, Tuple2<String, Long>> kvWordsDS = wordsDS
                .map(word -> Tuple2.of(word, 1L))
                .returns(Types.TUPLE(Types.STRING, Types.LONG));

        // 6. 对单词进行分组
        kvWordsDS.groupBy(0).sum(1).print();


    }
}

在这里插入图片描述

流式处理

public class StreamWordCount {
    public static void main(String[] args) throws Exception {
        // 1. 准备环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        // 在DataStream中设置batch处理模式
//        env.setRuntimeMode(RuntimeExecutionMode.BATCH);
        // 2. 实时读取数据
        DataStreamSource<String> lineDS = env.readTextFile("./data/words.txt");

        // 3. 对数据进行切分操作
        SingleOutputStreamOperator<String> wordsDS = lineDS.flatMap((FlatMapFunction<String, String>) (line, out) -> {
            for (String word : line.split(" ")) {
                out.collect(word);
            }
        }).returns(Types.STRING);

        // 4. 对单词进行计数
        SingleOutputStreamOperator<Tuple2<String, Long>> kvWordsDS = wordsDS
                .map(word -> Tuple2.of(word, 1L))
                .returns(Types.TUPLE(Types.STRING, Types.LONG));

        // 5. 分区,聚合统计
        KeyedStream<Tuple2<String, Long>, String> value = kvWordsDS
                .keyBy( value1 -> value1.f0);
        value.sum(1).print();

        // 6. 执行任务
        env.execute();

    }
}

在这里插入图片描述

更多推荐