• Flink Java Table API & SQL 之 wordcount


    Table API

    package com.daidai.table;
    
    import org.apache.flink.api.common.functions.MapFunction;
    import org.apache.flink.api.common.typeinfo.TypeInformation;
    import org.apache.flink.api.common.typeinfo.Types;
    import org.apache.flink.api.java.tuple.Tuple2;
    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.table.api.Table;
    import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
    import org.apache.flink.types.Row;
    
    import java.util.Arrays;
    
    import static org.apache.flink.table.api.Expressions.$;
    
    public class WordCountTable {
        public static void main(String[] args) throws Exception {
            StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
            StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
    
            DataStreamSource<String> source = env.fromCollection(Arrays.asList("hello", "word", "java", "scala", "java"));
            SingleOutputStreamOperator<Tuple2<String, Integer>> wordAndOne = source.map(new MapFunction<String, Tuple2<String, Integer>>() {
                @Override
                public Tuple2<String, Integer> map(String value) throws Exception {
                    return Tuple2.of(value, 1);
                }
            });
    
            Table table = tableEnv.fromDataStream(wordAndOne, $("word"), $("sum"));
            Table result = table.groupBy($("word"))
                    .select($("word"), $("sum").sum().as("count"));
    
            tableEnv.toRetractStream(result, Row.class).print();
    
            env.execute();
        }
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39

    SQL

    package com.daidai.table;
    
    import org.apache.flink.api.common.functions.MapFunction;
    import org.apache.flink.api.java.tuple.Tuple2;
    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.table.api.Table;
    import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
    import org.apache.flink.types.Row;
    
    import java.util.Arrays;
    
    import static org.apache.flink.table.api.Expressions.$;
    
    public class WordCountSQL {
        public static void main(String[] args) throws Exception {
            StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
            StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
    
            DataStreamSource<String> source = env.fromCollection(Arrays.asList("hello", "word", "java", "scala", "java"));
            SingleOutputStreamOperator<Tuple2<String, Integer>> wordAndOne = source.map(new MapFunction<String, Tuple2<String, Integer>>() {
                @Override
                public Tuple2<String, Integer> map(String value) throws Exception {
                    return Tuple2.of(value, 1);
                }
            });
    
            Table table = tableEnv.fromDataStream(wordAndOne, $("word"), $("sum"));
    
            tableEnv.createTemporaryView("wc", table);
    
            Table result = tableEnv.sqlQuery("select word, sum(`sum`) from wc group by word");
    
            tableEnv.toRetractStream(result, Row.class).print();
    
            env.execute();
        }
    
    }
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24
    • 25
    • 26
    • 27
    • 28
    • 29
    • 30
    • 31
    • 32
    • 33
    • 34
    • 35
    • 36
    • 37
    • 38
    • 39
    • 40
  • 相关阅读:
    Web系统常见安全漏洞介绍及解决方案-CSRF攻击
    CCF ChinaSoft 2023 论坛巡礼|形式验证@EDA论坛
    C#:变量的更多内容
    2021-09-29破解小米“铁蛋”,只需9999元,你也可以做一个四足机器人!
    一台服务器通过nginx安装多个web应用
    Day07--wxs的概念以及其基本的用法
    Vue.js 框架源码与进阶 - Vue.js 源码剖析 -虚拟 DOM
    Spring Boot Actuator 模块,spring-boot-starter-actuator
    Django+Vue.js学习记录(一)环境安装与配置
    NoSQL数据库入门
  • 原文地址:https://blog.csdn.net/weixin_46376562/article/details/125635530