Appearance
大数据系统-函数式编程
- Scala
- Spark RDD
- Flink 算子
- Java Stream API
很多人学 Spark、Flink,一上来就背:
java
map
flatMap
filter
reduce
join
window结果背完发现:
根本不知道这些算子到底在干什么。
其实用费曼方法讲,所有大数据框架(Spark/Flink/Beam)的算子,本质只有4类:
text
1. 变形
2. 过滤
3. 聚合
4. 关联剩下90%的算子都是这四类的变种。
先理解RDD是什么
RDD 是 Spark 最经典的抽象。
全称:
text
Resilient Distributed Dataset
弹性分布式数据集可以理解成:
text
超大号List例如:
java
List<String> data =
[
"张三",
"李四",
"王五"
]在 Spark 里:
java
RDD<String>只是这个 List:
text
存储在很多机器上例如:
text
机器1
张三
机器2
李四
机器3
王五一切算子都在操作RDD
例如:
java
rdd
.map(...)
.filter(...)
.reduce(...)其实就是:
text
RDD
↓
RDD
↓
RDD
↓
结果流水线加工。
第一类:Map(变形)
最重要的算子。
map
作用:
text
一条变一条例如:
java
[1,2,3]变平方:
java
rdd.map(x -> x*x)结果:
java
[1,4,9]费曼理解:
你有一筐苹果:
text
苹果
苹果
苹果map就是:
text
削皮每个苹果变样子。
但数量不变。
text
苹果3个
↓
削皮苹果3个第二类:FlatMap(爆裂)
很多新人最容易懵。
输入:
java
[
"hello world",
"flink spark"
]map:
java
map(str -> str.split(" "))结果:
java
[
["hello","world"],
["flink","spark"]
]flatMap:
java
flatMap(str -> str.split(" "))结果:
java
[
"hello",
"world",
"flink",
"spark"
]费曼理解:
map:
text
箱子变箱子flatMap:
text
拆箱例如:
text
快递箱1
苹果
香蕉
快递箱2
葡萄flatMap:
text
苹果
香蕉
葡萄全部倒出来。
第三类:Filter(筛子)
作用:
text
留下满足条件的数据例如:
java
[1,2,3,4,5]保留偶数:
java
filter(x -> x%2==0)结果:
java
[2,4]费曼理解:
像筛沙子。
text
石头
沙子
石头
沙子过滤:
text
只保留沙子第四类:Reduce(合并)
作用:
text
多条变一条例如:
java
[1,2,3,4]求和:
java
reduce((a,b)->a+b)过程:
text
1+2=3
3+3=6
6+4=10结果:
text
10费曼理解:
像滚雪球。
text
小雪球
↓
中雪球
↓
大雪球
↓
超级雪球GroupBy(分组)
SQL:
sql
GROUP BY对应:
java
groupBy例如:
text
张三 男
李四 男
小红 女分组:
text
男
├ 张三
└ 李四
女
└ 小红费曼理解:
学生按班级站队。
KeyBy(Flink最核心)
Spark:
java
groupByKeyFlink:
java
keyBy例如:
订单:
text
A用户
A用户
B用户
A用户keyBy:
java
keyBy(userId)变成:
text
A -> Task1
B -> Task2作用:
text
相同Key进入同一个状态空间这是状态计算的基础。
Aggregate(聚合)
例如:
java
sum
max
min
avg输入:
text
100
200
300sum:
text
600费曼理解:
会计记账。
Join(关联)
SQL:
sql
join用户表:
text
1 张三
2 李四订单表:
text
1 手机
2 电脑join:
text
张三 手机
李四 电脑费曼理解:
相亲。
text
男生名单
+
女生名单按照条件配对。
Window(窗口)
Flink灵魂算子。
因为流数据:
text
一直来
一直来
一直来没有结束。
怎么统计?
切时间片。
text
10:00~10:05
10:05~10:10
10:10~10:15这就是Window。
例如:
java
window(Tumbling 5min)统计:
text
5分钟订单数费曼理解:
彩票开奖。
text
每5分钟开奖一次统计这一期的数据。
Union(合流)
java
stream1.union(stream2)结果:
text
流1
+
流2合并。
费曼理解:
两条河汇成一条河。
Distinct(去重)
输入:
text
A
A
B
C
C输出:
text
A
B
C费曼理解:
签到表去重。
Spark与Flink最常用算子
Spark RDD
java
map
flatMap
filter
reduce
groupByKey
reduceByKey
join
union
distinct
sortBySpark Dataset/DataFrame
sql
select
where
group by
join
order by本质还是那些东西。
Flink DataStream
java
map
flatMap
filter
keyBy
window
process
aggregate
reduce
join
connect
union
sinkFlink独有大杀器:ProcessFunction
这是 Flink 的「九阴真经」。
前面那些:
java
map
filter
window都是封装好的套路。
ProcessFunction:
java
process()允许你:
text
操作状态
注册定时器
控制事件时间
处理迟到数据
自定义逻辑很多复杂业务:
text
实时风控
实时推荐
实时监控
CEP最后都会落到:
java
KeyedProcessFunction上。
一句话总结:
text
Map = 变形
FlatMap = 拆箱
Filter = 筛选
Reduce = 合并
GroupBy = 分组
KeyBy = 分桶
Aggregate= 汇总
Join = 配对
Window = 时间切片
Union = 合流
Distinct = 去重
Process = 自定义大招你如果把这些算子映射到 Java Stream:
java
stream()
.map()
.filter()
.collect()再映射到 SQL:
sql
SELECT
WHERE
GROUP BY
JOIN就会发现一个有意思的事实:
Spark、Flink、Java Stream、SQL 其实是在干同一件事——把一堆数据经过一条流水线,不断变形、筛选、聚合,最终产出结果。区别只是 Spark 处理分布式批数据,Flink 处理分布式流数据,而 SQL 处理数据库里的表。