flink中union 和connect学习记录

一、首先先配置flink的依赖,这里只用了flink-client客户端的类,其实包含了flink-java、flink-stream-java等依赖,不用再去单独引用

<!-- https://mvnrepository.com/artifact/org.apache.flink/flink-clients -->
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-clients_2.12</artifactId>
    <version>1.13.6</version>
</dependency>
<!-- https://mvnrepository.com/artifact/org.projectlombok/lombok -->
<dependency>
    <groupId>org.projectlombok</groupId>
    <artifactId>lombok</artifactId>
    <version>1.18.20</version>
</dependency>

flink-client包的依赖如下,包含了flink-java、flink-stream-java等依赖

 

 

二、union 详情

1、flink 中的union和sql中union all比较类似,类似从2张表张组成相向的查询字段进行union all合并成一张表;fink union 两个相同结构的stream合并相同结构的流

2、定义UserInfo类

@Data
public class UserInfo {
    private String id;
    private String name;
    private String gender;
    private Integer age;
    private Integer timestamp;

    public UserInfo() {
    }

    public UserInfo(String id, String name, String gender, Integer age, Integer timestamp) {
        this.id = id;
        this.name = name;
        this.gender = gender;
        this.age = age;
        this.timestamp = timestamp;
    }
}

 3、flink 程序

public class UnionTest {
    public static void main(String[] args) throws Exception {
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);
        // 处理java bean 成 UserInfo
        SingleOutputStreamOperator<UserInfo> stream1 = env.fromElements(
                new UserInfo("1001", "小张", "男", 21, 1),
                new UserInfo("1002", "小红", "女", 18, 2)
        );
        // 处理java bean 成 UserInfo
        SingleOutputStreamOperator<UserInfo> stream2 = env.fromElements(
                new UserInfo("2001", "老张", "男", 50, 4),
                new UserInfo("2002", "老红", "女", 53, 5),
                new UserInfo("2003", "老王", "男", 54, 6)
        );

        // 使用union
        stream1.union(stream2)
                .process(new ProcessFunction<UserInfo, String>() {
                    @Override
                    public void processElement(UserInfo value, Context ctx, Collector<String> out) throws Exception {
                        out.collect(value.toString());
                    }
                })
                .print();

        // 启动程序
        env.execute();
    }
}

4、 结果:

 

三、connect详情

connect有两张使用方式:

1、第一种和union 一样,多少条数据进,多少条数据出;

// Person 类数据定义
SingleOutputStreamOperator<Person> stream1 = env.fromElements(
        new Person("1002", 21, 2),
        new Person("1003", 32, 3),
        new Person("1004", 41, 4)
);

// PersonDetail 类数据定义
SingleOutputStreamOperator<PersonDetail> stream2 = env.fromElements(
        new PersonDetail("1001", "zhangsan", "man", 2),
        new PersonDetail("1002", "lisi", "woman", 3),
        new PersonDetail("1003", "wangwu", "man", 4),
        new PersonDetail("1004", "shugj", "man", 5)
);

// 第一种,类似union all ,多少条数据进,多少数据出
// 不分区使用
stream1.connect(stream2)
        .process(new CoProcessFunction<Person, PersonDetail, Tuple2<String, String>>() {
            @Override
            public void processElement1(Person value, Context ctx, Collector<Tuple2<String, String>> out) throws Exception {
                // 这里把id,age硬拼进去去,是否为了展示
                out.collect(new Tuple2<>(value.getId(),value.getAge().toString()));
            }

            @Override
            public void processElement2(PersonDetail value, Context ctx, Collector<Tuple2<String, String>> out) throws Exception {
                // 这里把id,name硬拼进去去,是否为了展示
                out.collect(new Tuple2<>(value.getId(),value.getName()));
            }
        }).print();
env.execute();

输出结果:

 

 

2、 第二种和join一样,只和等值关联上的输出;

// Person 类数据定义
        SingleOutputStreamOperator<Person> stream1 = env.fromElements(
                new Person("1002", 21, 2),
                new Person("1003", 32, 3),
                new Person("1004", 41, 4)
        );

        // PersonDetail 类数据定义
        SingleOutputStreamOperator<PersonDetail> stream2 = env.fromElements(
                new PersonDetail("1001", "zhangsan", "man", 2),
                new PersonDetail("1002", "lisi", "woman", 3),
                new PersonDetail("1003", "wangwu", "man", 4),
                new PersonDetail("1004", "shugj", "man", 5)
        );

// 第二种,这种类似join,等值关联多少条数据
// 分区
stream1.keyBy(data -> data.getId()) // 分区
        .connect(stream2.keyBy(data -> data.getId()))
        .process(new KeyedCoProcessFunction<String, Person, PersonDetail, Tuple2<Person,PersonDetail>>() {
            ValueState<Person> personState;
            ValueState<PersonDetail> personDetailState;

            @Override
            public void open(Configuration parameters) throws Exception {
                personState = getRuntimeContext().getState(new ValueStateDescriptor<Person>("person",Person.class));
                personDetailState = getRuntimeContext().getState(new ValueStateDescriptor<PersonDetail>("person-detail",PersonDetail.class));
            }

            @Override
            public void onTimer(long timestamp, OnTimerContext ctx, Collector<Tuple2<Person, PersonDetail>> out) throws Exception {
                personState.clear();
                personDetailState.clear();
            }

            @Override
            public void processElement1(Person value, Context ctx, Collector<Tuple2<Person, PersonDetail>> out) throws Exception {
                // 获取对方的状态数据
                PersonDetail personDetail = personDetailState.value();

                if (personDetail != null) {
                    // 如果 personDetail 有数据,就输出数据,并清空状态
                    out.collect(new Tuple2<>(value, personDetail));
                    personState.clear();
                    personDetailState.clear();
                } else {
                    // 如果 personDetail,没有数据,就把 person写入状态数据
                    personState.update(value);

                    // 开启定时器,10分钟后对方状态还没来就清空状态
                    // 获取当前时间
                    ctx.timerService().registerProcessingTimeTimer(value.getTimestamp() + 10*60*1000L);
                }
            }

            @Override
            public void processElement2(PersonDetail value, Context ctx, Collector<Tuple2<Person, PersonDetail>> out) throws Exception {
                // 获取对方的状态数据
                Person person = personState.value();

                if (person != null) {
                    // 如果 person 有数据,就输出数据,并清空状态
                    out.collect(new Tuple2<>(person, value));
                    personState.clear();
                    personDetailState.clear();
                } else {
                    // 如果没有数据,就把personDetail写入状态
                    personDetailState.update(value);

                    // 开启定时器,10分钟后对方状态还没来就清空状态
                    // 获取当前时间
                    ctx.timerService().registerProcessingTimeTimer(value.getTimestamp() + 10*60*1000L);
                }
            }
        }).print();
env.execute();

输出结果:

 

 

版权声明:本文内容由互联网用户自发贡献,该文观点与技术仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌侵权/违法违规的内容, 请发送邮件至 dio@foxmail.com 举报,一经查实,本站将立刻删除。

相关推荐


学习编程是顺着互联网的发展潮流,是一件好事。新手如何学习编程?其实不难,不过在学习编程之前你得先了解你的目的是什么?这个很重要,因为目的决定你的发展方向、决定你的发展速度。
IT行业是什么工作做什么?IT行业的工作有:产品策划类、页面设计类、前端与移动、开发与测试、营销推广类、数据运营类、运营维护类、游戏相关类等,根据不同的分类下面有细分了不同的岗位。
女生学Java好就业吗?女生适合学Java编程吗?目前有不少女生学习Java开发,但要结合自身的情况,先了解自己适不适合去学习Java,不要盲目的选择不适合自己的Java培训班进行学习。只要肯下功夫钻研,多看、多想、多练
Can’t connect to local MySQL server through socket \'/var/lib/mysql/mysql.sock问题 1.进入mysql路径
oracle基本命令 一、登录操作 1.管理员登录 # 管理员登录 sqlplus / as sysdba 2.普通用户登录
一、背景 因为项目中需要通北京网络,所以需要连vpn,但是服务器有时候会断掉,所以写个shell脚本每五分钟去判断是否连接,于是就有下面的shell脚本。
BETWEEN 操作符选取介于两个值之间的数据范围内的值。这些值可以是数值、文本或者日期。
假如你已经使用过苹果开发者中心上架app,你肯定知道在苹果开发者中心的web界面,无法直接提交ipa文件,而是需要使用第三方工具,将ipa文件上传到构建版本,开...
下面的 SQL 语句指定了两个别名,一个是 name 列的别名,一个是 country 列的别名。**提示:**如果列名称包含空格,要求使用双引号或方括号:
在使用H5混合开发的app打包后,需要将ipa文件上传到appstore进行发布,就需要去苹果开发者中心进行发布。​
+----+--------------+---------------------------+-------+---------+
数组的声明并不是声明一个个单独的变量,比如 number0、number1、...、number99,而是声明一个数组变量,比如 numbers,然后使用 nu...
第一步:到appuploader官网下载辅助工具和iCloud驱动,使用前面创建的AppID登录。
如需删除表中的列,请使用下面的语法(请注意,某些数据库系统不允许这种在数据库表中删除列的方式):
前不久在制作win11pe,制作了一版,1.26GB,太大了,不满意,想再裁剪下,发现这次dism mount正常,commit或discard巨慢,以前都很快...
赛门铁克各个版本概览:https://knowledge.broadcom.com/external/article?legacyId=tech163829
实测Python 3.6.6用pip 21.3.1,再高就报错了,Python 3.10.7用pip 22.3.1是可以的
Broadcom Corporation (博通公司,股票代号AVGO)是全球领先的有线和无线通信半导体公司。其产品实现向家庭、 办公室和移动环境以及在这些环境...
发现个问题,server2016上安装了c4d这些版本,低版本的正常显示窗格,但红色圈出的高版本c4d打开后不显示窗格,
TAT:https://cloud.tencent.com/document/product/1340