flink的一些笔记
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








更多推荐
所有评论(0)