1、任务槽slot

slot特点: 均分隔离内存,不隔离cpu

                        可以共享:同一个job中,不同算子的子任务才可以共享 同一个slot,并且是同时在运行的,前提是属于同一个slot共享组,默认都是“default”

slot数量与并行度的关系

        slot是一种静态的概念,表示最大的并发上限

        并行度是一种动态的概念,表示实际运行占用了几个

        要求:slot数量>=job并发度(算子最大并行度),job才能运行

        注意:如果是yarn模式,动态申请

                --》申请的tm数量=job并行度 / 每个tm的slot数,向上取整

        比如session:一开始0个taskmanager,0个slot

                --》提交一个job(每个tm的slot数为3),并行度10

                        --》10/3,向上取整,申请4个tm

2、转换算子transformation

map        一一对应

filter        true保留,false去掉

flatmap        一进多出(一进一出、一进零出(就相当于过滤)、一进多出)

map怎么控制一进一出   =》 使用return

flatmap怎么控制一进多出 =》 通过collector来输出,调用几次就输出几条

简单聚合算子

有了按键分区的数据流keyedstream(必须经过keyby之后才会有),就可以基于它进行聚合操作

sum()、min()、max()、

minBy()与min()类似,在输入流上针对指定字段求最小值。不同的是,min()只计算指定字段的最小值,其他字段会保留最初第一个数据的值,而minBy()则会返回包含字段最小值的整条数据

规约聚合reduce

同一个 key,每来一条新数据,就拿【历史聚合结果】和【新来的数据】做自定义计算,输出最新聚合结果。

用户自定义函数

用户自定义函数分为;函数类、匿名函数、富函数类

分区算子

分流

filter分:类比把数据拆分为奇数流和偶数流,每个流需要单独判断

侧输出流

就是把数据比如s1,s2,s3,s4,s5.....拆分为支流s1,支流s2,剩下的都是主流

合流

用connect,有两个map,2条流进来之后各处理各的

3、keyby     

 对于flink,datastream是没有直接进行聚合的api的,所以做聚合前要先分区并行处理,keyby通过指定键key,可以将一条流从逻辑上划分成不同的分区(也就是并行处理的子任务),基于不同的key,流中的数据将被分配到不同的分区,所有相同key的数据发到同一个分区

同一个分区内可以有多个组

4、输出算子sink

输出到文件

输出到kafka

输出到mysql

自定义sink输出(不推荐)

5、窗口

类比为桶是现找的,前一个桶满了之后才去找下一个桶,不是一次性把需要多少桶一次性找齐

基于时间、个数

基于时间、个数

会话窗口重点在于时间间隔,基于时间

6、时间语义

7、水位线

8、基于时间的双流联结

9、处理函数

10、容错机制

端到端一致性

11、flinksql

更多推荐