ICode9

精准搜索请尝试: 精确搜索
首页 > 其他分享> 文章详细

Flink11--FliterAndKeyBy算子

2022-03-27 17:32:16  阅读:148  来源: 互联网

标签:Flink11 -- flink videoOrder api FliterAndKeyBy org apache import


一、导入依赖

参考本人下博客

二、代码

FLink11FilterApp.java

package net.xdclass.class9;

import org.apache.flink.api.common.RuntimeExecutionMode;
import org.apache.flink.api.common.functions.FilterFunction;
import org.apache.flink.api.java.functions.KeySelector;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

import net.xdclass.source.VideoOrderSourceV2;
import net.xdclass.model.VideoOrder;

/**
 * @desc filter算子
 * @menu
 */
public class FLink11FilterApp {

    public static void main(String[] args) throws Exception{
        //WebUi方式运行
//        final StreamExecutionEnvironment env =
//                StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(new Configuration());
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        //设置运行模式为流批一体
        env.setRuntimeMode(RuntimeExecutionMode.AUTOMATIC);
        //并行度
        env.setParallelism(1);
        //设置为自定义source
        DataStream<VideoOrder> ds = env.addSource(new VideoOrderSourceV2());

        SingleOutputStreamOperator<VideoOrder> filterDs = ds.filter(new FilterFunction<VideoOrder>() {
            @Override
            public boolean filter(VideoOrder videoOrder) throws Exception {
                return videoOrder.getMoney() > 10;
            }
        });

//        KeyedStream<VideoOrder, Object> videoKeyBy = filterDs.keyBy(new KeySelector<VideoOrder, Object>() {
//            @Override
//            public Object getKey(VideoOrder videoOrder) throws Exception {
//                return videoOrder.getTitle();
//            }
//        });
//        SingleOutputStreamOperator<VideoOrder> videoKeySum = videoKeyBy.sum("money");
        SingleOutputStreamOperator<VideoOrder> moneyDs = filterDs.keyBy(new KeySelector<VideoOrder, Object>() {
            @Override
            public Object getKey(VideoOrder videoOrder) throws Exception {
                return videoOrder.getTitle();
            }
        }).sum("money");

        moneyDs.print();

        //DataStream需要调用execute,可以取个名称
        env.execute("keyBy map job");
    }
}

 

标签:Flink11,--,flink,videoOrder,api,FliterAndKeyBy,org,apache,import
来源: https://www.cnblogs.com/robots2/p/16063576.html

本站声明: 1. iCode9 技术分享网(下文简称本站)提供的所有内容,仅供技术学习、探讨和分享;
2. 关于本站的所有留言、评论、转载及引用,纯属内容发起人的个人观点,与本站观点和立场无关;
3. 关于本站的所有言论和文字,纯属内容发起人的个人观点,与本站观点和立场无关;
4. 本站文章均是网友提供,不完全保证技术分享内容的完整性、准确性、时效性、风险性和版权归属;如您发现该文章侵犯了您的权益,可联系我们第一时间进行删除;
5. 本站为非盈利性的个人网站,所有内容不会用来进行牟利,也不会利用任何形式的广告来间接获益,纯粹是为了广大技术爱好者提供技术内容和技术思想的分享性交流网站。

专注分享技术,共同学习,共同进步。侵权联系[81616952@qq.com]

Copyright (C)ICode9.com, All Rights Reserved.

ICode9版权所有