Skip to content

大数据系统-函数式编程 ​

  • 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
groupByKey

Flink:

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
300

sum:

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
sortBy

Spark Dataset/DataFrame ​

sql
select
where
group by
join
order by

本质还是那些东西。


java
map
flatMap
filter
keyBy
window
process
aggregate
reduce
join
connect
union
sink

Flink独有大杀器: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 处理数据库里的表。