Flink流式计算WordCountTopN

Flink流式计算WordCountTopN可以采用流处理编程和FlinkSql自定义UDTF函数的方式

流处理编程方法:

public class Flink05_WC_TOPN {

public static void main(String[] args)throws Exception {

//TODO 1.获取执行环境

        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

      // StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(new Configuration());

//并行度1

        env.setParallelism(1);

        //读取无界数据

        DataStreamSource streamSource = env.socketTextStream("hadoop102", 8888);

        //TODO 2.切分数据

        SingleOutputStreamOperator outputStreamOperator = streamSource.flatMap(new FlatMapFunction() {

@Override

            public void flatMap(String value, Collector out)throws Exception {

String[] strings = value.split(" ");

                for (String string : strings) {

out.collect(string);

                }

}

});

        //转为tuple

        SingleOutputStreamOperator> streamOperator = outputStreamOperator.map(new MapFunction>() {

@Override

            public Tuple2map(String value)throws Exception {

return Tuple2.of(value, 1);

            }

});

        ArrayList> top3list =new ArrayList<>();

        //TODO 3. 聚合求sum

        KeyedStream, Tuple> tuple2TupleKeyedStream = streamOperator.keyBy(0);

        SingleOutputStreamOperator> KVstream = tuple2TupleKeyedStream.sum(1);

        SingleOutputStreamOperator TOP3 = KVstream.process(new ProcessFunction, String>() {

@Override

            //除去此String在List中保存的上一个状态

            public void processElement(Tuple2 value, Context ctx, Collector out)throws Exception {

for (int i =0; i

Tuple2 stringIntegerTuple2 =top3list.get(i);

                    if(stringIntegerTuple2.f0.equals(value.f0)){

top3list.remove(i);

                    }

}

top3list.add(value);

                top3list.sort(new Comparator>() {

@Override

                    public int compare(Tuple2 o1, Tuple2 o2) {

return Integer.valueOf(o2.f1) - Integer.valueOf(o1.f1);

                    }

});

                if (top3list.size() >3) {

top3list.remove(3);

                }

out.collect(top3list.toString());

            }

});

        TOP3.print();

        env.execute();

    }

}


FlinkSql自定义UDATF函数的方式使用

public class Flink26_Fun_UDATF {

public static void main(String[] args) {

//TODO 1.获取流执行环境

        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        env.setParallelism(1);

        //TODO 2.获取表的执行环境

        StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);

        //TODO 3.获取数据

        DataStreamSource streamSource = env.fromElements(new WaterSensor("sensor_1", 1000L, 10),

                new WaterSensor("sensor_1", 2000L, 20),

                new WaterSensor("sensor_2", 3000L, 30),

                new WaterSensor("sensor_1", 4000L, 40),

                new WaterSensor("sensor_1", 5000L, 50),

                new WaterSensor("sensor_2", 6000L, 60));

        //2.从端口读取数据

/*      DataStreamSource Source= env.socketTextStream("hadoop102", 8888);

//3.将数据转为WaterSensor

SingleOutputStreamOperator streamSource = Source.map(new 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]));

}

});

*/

        //TODO 4.将流转为表

        Table table = tableEnv.fromDataStream(streamSource, $("id"), $("ts"), $("vc"));

        //TODO 5.使用UDF

        //不注册直接使用

        table

.groupBy("id")

.flatAggregate(call(MyUDATF.class,$("vc")).as("value","top"))

.select($("id"),$("value"),$("top"))

//  .select($("id"),$("f0"),$("f1"))  //不as calue 和top写法

                .execute().print();

        //注册一个定义函数

/*  tableEnv.createTemporarySystemFunction("MyUDATF", MyUDATF.class);

//TableAPi

table

.groupBy("id")

.flatAggregate(call("MyUDATF",$("vc")).as("value","top"))

.select($("id"),$("value"),$("top"))

.execute().print();*/

    }

//定义一个类当做累加器

    public static class top2Vc{

public Integerfirst=Integer.MIN_VALUE;;

        public Integersecond = Integer.MIN_VALUE;

    }

//自定义UDATF函数,实现TOP2

  public  static class MyUDATFextends TableAggregateFunction,top2Vc>{

//创建累加器

        @Override

        public top2VccreateAccumulator() {

top2Vc top2Vc =new top2Vc();

            return top2Vc;

        }

//累加操作

        public  void accumulate(top2Vc acc,Integer value){

//先比较当前数据是否大于第一

            if(value>acc.first){

acc.second=acc.first;

                acc.first=value;

            }else if(value>acc.second) {

acc.second=value;

            }

}

//将结果发送出去

        public void emitValue(top2Vc acc, Collector> out){

if(acc.first!=Integer.MIN_VALUE){

out.collect(Tuple2.of(acc.first,1));

            }if(acc.second!=Integer.MIN_VALUE){

out.collect(Tuple2.of(acc.second,2));

            }

}

}

}

©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 218,546评论 6 507
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 93,224评论 3 395
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 164,911评论 0 354
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 58,737评论 1 294
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 67,753评论 6 392
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 51,598评论 1 305
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 40,338评论 3 418
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 39,249评论 0 276
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 45,696评论 1 314
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 37,888评论 3 336
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 40,013评论 1 348
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 35,731评论 5 346
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 41,348评论 3 330
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 31,929评论 0 22
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 33,048评论 1 270
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 48,203评论 3 370
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 44,960评论 2 355

推荐阅读更多精彩内容