0%

Flink-SQL-active

维表同步脚本:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
insert
into
dim_result_lct_activy_config
select
Fact_id,
LAST_VALUE(Fact_name),
FROM_UNIXTIME(UNIX_TIMESTAMP(LAST_VALUE(Fact_start_time)),
'yyyyMMddHHmmss') as startTime,
FROM_UNIXTIME(UNIX_TIMESTAMP(LAST_VALUE(Fact_end_time)),
'yyyyMMddHHmmss') as endTime,
LAST_VALUE(Fstate)
from
db_act_config_t_act_logic_config
group by
Fact_id

FlinkSQL Upsert/Retraction 写入 MySQL 的问题

retract流是算子的特性,不同的算子不同,比如group by聚合算子天生就是会发retract流的,比如你之前算出来的是a, 1, 后面变化了,变成了a,2,所以group by算子会发-a,1, +a,2.

像join的话也是有retract流的,比如outer join

当然如果下游sink支持upsert, 这块也会有优化,会优化为update。

INSERT INTO mysql_sink SELECT f1, count(*) FROM kafka_src GROUP BY f1
每从 kafka 过来一条新的记录,会生成两条记录 Tuple2<Row, Boolean>, 旧的被删除,新的会添加上。这是query是会一个会产生retract stream的query,可以简单理解成每条kafka的数据过来会产生两条记录,但是最终写入下游的系统。需要看下游的系统支持和实现的sink(现在有三种sink AppendStreamSink, UpsertStreamSink, RetractStreamSink)

我看 https://github.com/apache/flink/tree/master/flink-connectors/flink-jdbc/src/main/java/org/apache/flink/api/java/io/jdbc 没有 Retract 方式
实际上使用了 JDBCUpsertTableSink.java 的代码写入 MySQL 吗?
现有的sink中,kafka是实现的AppendStreamSink,所以只支持insert 的记录,不支持retract.
你用DDL声明的mysql表,对应的jdbc sink 是JDBCUpsertTableSink,所以会按照Upsert的逻辑处理, 也不支持retract。

如若不带 group by 直接:
INSERT INTO mysql_sink SELECT f1, f2 FROM kafka_src
主键冲突写入 mysql 是会出错的,怎么可以用 Upsert 的方式直接覆盖呢?

不带 group by时无法推导出query的 unique key,没法做按照unique key的更新,
只需要将 query的 key (你这里是group by 后的字段)和db中主键保持一致即可

TopN

1
2
3
4
5
6
7
SELECT [column_list]
FROM (
SELECT [column_list],
ROW_NUMBER() OVER ([PARTITION BY col1[, col2...]]
ORDER BY col1 [asc|desc][, col2 [asc|desc]...]) AS rownum
FROM table_name)
WHERE rownum <= N [AND conditions]
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
EXPLAIN PLAN FOR SELECT * FROM (
SELECT *,
row_number() OVER(PARTITION BY merchandiseId ORDER BY totalQuantity DESC) AS rownum
FROM (
SELECT merchandiseId, sum(quantity) AS totalQuantity
FROM rtdw_dwd.kafka_order_done_log
GROUP BY merchandiseId
)
) WHERE rownum <= 10

== Abstract Syntax Tree ==
LogicalProject(merchandiseId=[$0], totalQuantity=[$1], rownum=[$2])
+- LogicalFilter(condition=[<=($2, 10)])
+- LogicalProject(merchandiseId=[$0], totalQuantity=[$1], rownum=[ROW_NUMBER() OVER (PARTITION BY $0 ORDER BY $1 DESC NULLS LAST)])
+- ...

== Optimized Logical Plan ==
Rank(strategy=[RetractStrategy], rankType=[ROW_NUMBER], rankRange=[rankStart=1, rankEnd=10], partitionBy=[merchandiseId], orderBy=[totalQuantity DESC], select=[merchandiseId, totalQuantity, w0$o0])
+- Exchange(distribution=[hash[merchandiseId]])
+- ...

== Physical Execution Plan ==
Stage 1 : Data Source
...

Stage 2 : Operator
...

Stage 4 : Operator
...

Stage 6 : Operator
content : Rank(strategy=[RetractStrategy], rankType=[ROW_NUMBER], rankRange=[rankStart=1, rankEnd=10], partitionBy=[merchandiseId], orderBy=[totalQuantity DESC], select=[merchandiseId, totalQuantity, w0$o0])
ship_strategy : HASH

滚动窗口 TUMBLE window

每分钟pv、uv

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
-- 滚动窗口统计pv:计算每分钟用户的页面访问数量
-- 上游表为tube,计算结果写入msql
-- 假设上游表为页面点击记录表,计算pv只需要使用count函数进行条数统计即可
INSERT INTO
console
SELECT
-- 滚动窗口的开始时间
FROM_UNIXTIME(
CAST(
TUMBLE_START(up.fevent_time, INTERVAL '1' MINUTE) as BIGINT
) / 1000,
'yyyyMMddHHmmss'
) as start_time,
count(1) as pv,
count(distinct fuin) as uv
FROM
datagen AS up
GROUP BY
-- 按照1分钟进行开窗统计
-- 第一个参数指定事件时间
-- 第二个参数指定多久触发一次统计
TUMBLE(up.fevent_time, INTERVAL '1' MINUTE) -- 查看结果
-- 打开TM日志,全局搜索ConSoleSink output,即刻查看每分钟统计的数据
1
2
3
4
5
6
7
8
9
10
11
12
INSERT INTO
t_hj_sample_result_040901
SELECT
CAST(up.Fact_id AS INTEGER) AS catalog_01,
CAST(
TUMBLE_END(up.ptime, INTERVAL '10' MINUTE) AS BIGINT
) AS window_end
FROM
b_cdg_cft_lct_act_data_db_act_t_user_prize AS up
GROUP BY
CAST(up.Fact_id AS INTEGER),
TUMBLE(up.ptime, INTERVAL '10' MINUTE)

滑动窗口

每小时统计一次过去一天的pv

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
-- 滑动窗口统计pv:每小时统计一次过去1天的用户页面访问数量
-- 上游表为tube,计算结果写入msql
-- 假设上游表为页面点击记录表,计算pv只需要使用count函数进行条数统计即可
INSERT INTO
console
SELECT
-- 滚动窗口的开始时间
FROM_UNIXTIME(
CAST(
HOP_START(
up.fevent_time,
INTERVAL '1' HOUR,
INTERVAL '1' DAY
) as BIGINT
) / 1000,
'yyyyMMddHHmmss'
) as start_time,
count(1) as pv
FROM
datagen AS up
GROUP BY
-- 按照1小时进行滑动开窗统计
-- 第一个参数指定事件时间
-- 第二个参数指定多久触发一次统计
-- 第三个参数指定统计多久之前的数据
HOP(
up.fevent_time,
INTERVAL '1' HOUR,
INTERVAL '1' DAY
) -- 查看结果
-- 打开TM日志,全局搜索ConSoleSink output,即刻查看每分钟统计的数据

每小时统计一次过去一周的各页面的pv

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
-- 滑动窗口统计uv:每天统计一次过去一周网页的独立用户访问量
-- 上游表为tube,写入mysql
-- 假设上游表为用户交易表,计算只需要对用户id进行去重即可
INSERT INTO
console
SELECT
-- 窗口的开始时间
FROM_UNIXTIME(
CAST(
HOP_START(up.fevent_time, INTERVAL '1' DAY, INTERVAL '7' DAY) as BIGINT
) / 1000,
'yyyyMMddHHmmss'
) as start_time,
-- 使用distinct算子进行去重
count(distinct(fuin)) as uv
From
datagen AS up
GROUP BY
-- 按照一天进行开窗统计
-- 第一个参数指定事件时间
-- 第二个参数指定多久触发一次统计
-- 第三个参数指定统计多久之前的数据
HOP(up.fevent_time, INTERVAL '1' DAY, INTERVAL '7' DAY) -- 查看结果
-- 打开TM日志,全局搜索ConSoleSink output,即刻查看每分钟统计的数据

增量窗口 INCREMENT

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
-- 增量窗口统计pv:每天统计一次用户页面访问数量,同时每分钟输出一次实时统计结果
-- 上游表为tube,计算结果写入msql
-- 假设上游表为页面点击记录表,计算pv只需要使用count函数进行条数统计即可
INSERT INTO
console
SELECT
-- 滚动窗口的开始时间
FROM_UNIXTIME(
CAST(
INCREMENT_START(
up.fevent_time,
INTERVAL '1' DAY,
INTERVAL '1' HOUR
) as BIGINT
) / 1000,
'yyyyMMddHHmmss'
) as start_time,
count(1) as pv
FROM
datagen AS up
GROUP BY
-- 按照1小时进行滑动开窗统计
-- 第一个参数指定事件时间
-- 第二个参数指定多久触发一次统计
-- 第三个参数指定多久输出一次实时结果
INCREMENT(
up.fevent_time,
INTERVAL '1' DAY,
INTERVAL '1' HOUR
) -- 查看结果
-- 打开TM日志,全局搜索ConSoleSink output,即刻查看每分钟统计的数据

1 Flink SQL

1.1 Fink SQL核心功能一览

  • SELECT FROM WHERE

  • GROUP BY /HAVING

  • Regular AGG

  • Time-windowed AGG(TUMBLE,HOP,SESSION)

  • Regular JOIN(INNER,LEFT/RIGHT/FULL)

  • Time-windowed JOIN(INNER,LEFT/RIGHT/FULL)

  • Temporal Table Join

  • 大量内置函数(150+)

  • LIKE, EXTRACT, TIMESTAMPADD, MD5, AVG..

  • 丰富的类型系统:POJO,Map,Array,Row,NestedType

  • 自定义函数(scalar,table,aggregate)

  • CEP on SQL(复杂事件处理)

https://flink.apache.org/2020/07/28/flink-sql-demo-building-e2e-streaming-application.html

http://wuchong.me/blog/2019/08/20/flink-sql-training/

https://github.com/ververica/flink-sql-cookbook

1.2 Retract mode 和 Append mode

toAppendStream 只支持insert
toRetractStream 其余模式都可以

如果动态表仅只有Insert操作,即之前输出的结果不会被更新,则使用该模式。如果更新或删除操作使用追加模式会失败报错,始终可以使用此模式。返回值是boolean类型。它用true或false来标记数据的插入和撤回,返回true代表数据插入,false代表数据的撤回。

使用flinkSQL处理实时数据当我们把表转化成流的时候,需要用toAppendStream与toRetractStream这两个方法。稍不注意可能直接选择了toAppendStream。

始终可以使用此模式。返回值是boolean类型。它用true或false来标记数据的插入和撤回,返回true代表数据插入,false代表数据的撤回。

当我们使用的sql语句包含:count() group by时,必须使用缩进模式

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
// 获取StreamTableEnvironment.
StreamTableEnvironment tableEnv = ...;
// 包含两个字段的表(String name, Integer age)
Table table = ...
// 将表转为DataStream,使用Append Mode追加模式,数据类型为Row
DataStream<Row> dsRow = tableEnv.toAppendStream(table, Row.class);
// 将表转为DataStream,使用Append Mode追加模式,数据类型为定义好的TypeInformation
TupleTypeInfo<Tuple2<String, Integer>> tupleType = new TupleTypeInfo<>(
Types.STRING(),
Types.INT());
DataStream<Tuple2<String, Integer>> dsTuple =
tableEnv.toAppendStream(table, tupleType);
// 将表转为DataStream,使用的模式为Retract Mode撤回模式,类型为Row
// 对于转换后的DataStream<Tuple2<Boolean, X>>,X表示流的数据类型,
// boolean值表示数据改变的类型,其中INSERT返回true,DELETE返回的是false
DataStream<Tuple2<Boolean, Row>> retractStream
tableEnv.toRetractStream(table, Row.class);

1.2.1 如何实现回退更新?

flink-connector-jdbc 最终使用的是SQL引擎的upsert语法:

1
insert into tbl() values ( ),( ) on duplicate key update

可以看下 MySQL/upsert 一节

对应flink 源码

1
org.apache.flink.connector.jdbc.dialect.MySQLDialect#getUpsertStatement

1.3 keyedBy与group by区别

1.4 SQL解析工具

hive使用了antlr3实现了自己的HQL,
Flink使用Apache Calcite,
而Calcite的解析器是使用JavaCC实现的,
Spark2.x以后采用了antlr4实现自己的解析器,
Presto也是使用antlr4。

1.5 通过Table api创建表

1
2
# create a Table from a Table API query
tapi_result = table_env.from_path("table1").select(...)

1.6 视图 view

1.7 window aggregate与group aggregate区别

参考自Apache Flink 零基础入门(九):Flink SQL 编程实践

1.7.1 Group Aggregate的例子

这是一个group aggregate, 内存中累计每个分类的数据,有变化时,update sink

1
2
3
4
SELECT psgCnt, COUNT(*) AS cnt 
FROM Rides
WHERE isInNYC(lon, lat)
GROUP BY psgCnt;

1.7.2 Window Aggregate

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
SELECT 
toAreaId(lon, lat) AS area, -- toAreaId为UDF
TUMBLE_END(rideTime, INTERVAL '5' MINUTE) AS window_end,
--滚动窗口的结束时间:
-- ① 可以时间不同吗?
-- ② 窗口可以别名吗?
COUNT(*) AS cnt
FROM Rides
WHERE isInNYC(lon, lat) and isStart
-- isStart 为boolean
GROUP BY
toAreaId(lon, lat),
TUMBLE(rideTime, INTERVAL '5' MINUTE)
-- 定义窗口
HAVING COUNT(*) >= 5;
-- 5分钟超过5次才输出

1.7.3 Window Aggregate 与 Group Aggregate 的区别

Window Aggregate 是当window结束时才输出,其输出的结果是最终值,不会再进行修改,其输出流是一个 Append 流。而 Group Aggregate 是每处理一条数据,就输出最新的结果,其结果是在不断更新的,就好像数据库中的数据一样,其输出流是一个 Update 流。

window 由于有 watermark ,可以精确知道哪些窗口已经过期了,所以可以及时清理过期状态,保证状态维持在稳定的大小。而 Group Aggregate 因为不知道哪些数据是过期的,所以状态会无限增长,这对于生产作业来说不是很稳定,所以建议对 Group Aggregate 的作业配上 State TTL 的配置。

例如统计每个店铺每天的实时PV,那么就可以将 TTL 配置成 24+ 小时,因为一天前的状态一般来说就用不到了。

1
2
3
SELECT  DATE_FORMAT(ts, 'yyyy-MM-dd'), shop_id, COUNT(*) as pv
FROM T
GROUP BY DATE_FORMAT(ts, 'yyyy-MM-dd'), shop_id

当然,如果 TTL 配置地太小,可能会清除掉一些有用的状态和数据,从而导致数据精确性地问题。这也是用户需要权衡地一个参数。

[参考文献]


  1. 在数据流中使用SQL查询:Apache Flink中的动态表的持续查询

  2. Flink Table API & SQL编程指南(1)

  3. Flink SQL金竹学习文档

Checkpoint API

算子支持checkpoint

CheckpointedFunction

org.apache.flink.streaming.api.checkpoint.CheckpointedFunction

CheckpointedFunction是Stateful transformation functions的核心接口,用于跨stream维护state:

public void snapshotState(FunctionSnapshotContext functionSnapshotContext) throws Exception
在checkpoint的时候会被调用,用于snapshot state,通常用于flush、commit、synchronize外部系统

public void initializeState(FunctionInitializationContext context) throws Exception
在parallel function初始化的时候(第一次初始化或者从前一次checkpoint recover的时候)被调用,通常用来初始化state,以及处理state recovery的逻辑

从checkpoint中恢复数据时,需要判断snapshot当前的情况,

FunctionSnapshotContext实现了ManagedSnapshotContext, 父类中的方法: getCheckpointId,getCheckpointTimestamp
FunctionInitializationContext实现了ManagedInitializationContext接口, 实现了 isRestoredgetOperatorStateStoregetKeyedStateStore方法

在初始化容器之后,我们使用上下文的isRestore()方法检查失败后是否正在恢复。如果是true,即正在恢复,则应用恢复逻辑。

样例: HBase写入OutPutFormat

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
/**
* WAL实现模式,这个模式有问题:
* 如果没有新数据进来,不会触发写入,造成下游看不到数据
*/
class PortraitOutputFormat extends RichOutputFormat<EventItem> implements CheckpointedFunction {
// 输出阈值,批量写入的条数
private final int threshold;
// 维护在状态中的数据
private transient ListState<EventItem> checkpointState;
// 内存中的数据
private List<EventItem> bufferedEventItem;
// HBase客户端
private HBaseClient hbaseClient;

public PortraitOutputFormat(HBaseClient hbaseClient) {
this.hbaseClient = hbaseClient;
}
/**
* checkpoint时调用
* 执行snapshot操作,将内存中的数据写入到内存
*/
@Override
public void snapshotState(FunctionSnapshotContext functionSnapshotContext) throws Exception {
checkpointState.clear();
for (EventItem eventItem : bufferedEventItem) {
checkpointState.add(eventItem);
}
}
/**
* 创建state,判断是否存在需要恢复的状态,如果有则需要恢复到bufferedEventItem
*/
@Override
public void initializeState(FunctionInitializationContext context) throws Exception {
ListStateDescriptor<EventItem> descriptor = new ListStateDescriptor<>("buf-p", EventItem.class);
checkpointState = context.getOperatorStateStore().getListState(descriptor);
if (context.isRestored()) {
for (EventItem eventItem : checkpointState.get()) {
bufferedEventItem.add(eventItem);
}
}
}
@Override
public void configure(Configuration configuration) {

}
@Override
public void open(int taskNumber, int numTasks) throws IOException {
}
/**
* 将新消息写入到缓存bufferedEventItem,缓存个数大约threshold,则执行sink写入,然后清空bufferedEventItem
*/
@Override
public void writeRecord(EventItem value) throws IOException {
if (value.getAttachUserId() == null) {
return;
}

bufferedEventItem.add(value);
int size = bufferedEventItem.size();

if (size >= threshold) {
List<Put> puts = bufferedEventItem
.stream()
.map(eventItem -> {
String rowKey1 = portraitDataGenerator.rowKey(eventItem);
Map<String, String> data = portraitDataGenerator.data(eventItem);
Put put = new Put(rowKey1.getBytes());
for (String cfc : data.keySet()) {
String[] cfcs = cfc.split(":");
String cf = cfcs[0];
String c = cfcs[1];
String dataOne = data.get(cfc);
put.addColumn(cf.getBytes(), c.getBytes(), dataOne.getBytes());
}
return put;
})
.collect(Collectors.toList());
try {
hbaseClient.putAndFlush(puts);
} catch (IOException e) {
e.printStackTrace();
}
bufferedEventItem.clear();
}
}

@Override
public void close() throws IOException {
if (hTable != null) {
hTable.flushCommits();
hTable.close();
}
if (connection != null) {
connection.close();
}
}
}

CheckpointListener

org.apache.flink.runtime.state.CheckpointListener

一旦所有checkpoint参与者确认完成,想要接收提交通知的功能/操作来实现。

TTL

1.8 自动清理原理

Flink 1.6.0版本引入了State TTL功能。它使流处理应用程序的开发人员配置过期时间,并在定义时间超时(Time to Live)之后进行清理。

在Flink 1.8.0中,该功能得到了扩展,包括对RocksDB和堆状态后端(FSStateBackend和MemoryStateBackend)的历史数据进行持续清理,从而实现旧条目的连续清理过程(根据TTL设置)。

RocksDB后台压缩可以过滤掉过期状态
如果你的Flink应用使用RocksDB作为状态后端存储,则可以启用另一个基于Flink特定压缩过滤器的清理策略。RocksDB定期运行异步压缩以合并状态更新并减少存储。Flink压缩过滤器使用TTL检查状态条目的到期时间戳,并丢弃所有过期值。

激活此功能的第一步是通过设置以下Flink配置选项来配置RocksDB状态后端:

state.backend.rocksdb.ttl.compaction.filter.enabled

配置RocksDB状态后端后,将为状态启用压缩清理策略,如以下代码示例所示:

1
2
3
4
StateTtlConfig ttlConfig = StateTtlConfig
.newBuilder(Time.days(7))
.cleanupInRocksdbCompactFilter()
.build();

TTL实现原理

State Backend 状态后端

三种状态后端:内存(MemoryStateend)、文件系统(FsStateend)和RocksDB(RocksDBStateend)

state 保存 snapshot与restore 大小
keyed state 堆内或堆外(RocksDB) backend自行实现,用户不关心
operator state 堆内 用户自行实现

图片

State backend snapshot保存 checkpoint保存
MemoryStateend 内存 内存
FsStateend 内存 文件系统,如hdfs
RocksDBStateend rocksdb 文件系统,如hdfs

Flink 的 keyed state 本质上来说就是一个键值对,所以与 RocksDB 的数据模型是吻合的。下图分别是 “window state” 和 “value state” 在 RocksDB 中的存储格式,所有存储的 key,value 均被序列化成 bytes 进行存储。

在 RocksDB 中,每个 state 独享一个 Column Family,而每个 Column family 使用各自独享的 write buffer 和 block cache,上图中的 window state 和 value state实际上分属不同的 column family。

最佳实践

operator state

慎重使用长 list

下图展示的是目前 task 端 operator state 在执行完 checkpoint 返回给 job master 端的 StateMetaInfo 的代码片段。

由于 operator state 没有 key group 的概念,所以为了实现改并发恢复的功能,需要对 operator state 中的每一个序列化后的元素存储一个位置偏移 offset,也就是构成了上图红框中的 offset 数组。

那么如果你的 operator state 中的 list 长度达到一定规模时,这个 offset 数组就可能会有几十 MB 的规模,关键这个数组是会返回给 job master,当 operator 的并发数目很大时,很容易触发 job master 的内存超用问题。我们遇到过用户把 operator state 当做黑名单存储,结果这个黑名单规模很大,导致一旦开始执行 checkpoint,job master 就会因为收到 task 发来的“巨大”的 offset 数组,而内存不断增长直到超用无法正常响应。

正确使用 UnionListState

union list state 目前被广泛使用在 kafka connector 中,不过可能用户日常开发中较少遇到,他的语义是从检查点恢复之后每个并发 task 内拿到的是原先所有operator 上的 state,如下图所示:

kafka connector 使用该功能,为的是从检查点恢复时,可以拿到之前的全局信息,如果用户需要使用该功能,需要切记恢复的 task 只取其中的一部分进行处理和用于下一次 snapshot,否则有可能随着作业不断的重启而导致 state 规模不断增长。

Keyed state 使用建议

如何正确清空当前的 state

state.clear() 实际上只能清理当前 key 对应的 value 值,如果想要清空整个 state,需要借助于 applyToAllKeys 方法,具体代码片段如下:

如果你的需求中只是对 state 有过期需求,借助于 state TTL 功能来清理会是一个性能更好的方案。

RocksDB 中考虑 value 值很大的极限场景

受限于 JNI bridge API 的限制,单个 value 只支持 2^31 bytes 大小,如果存在很极限的情况,可以考虑使用 MapState 来替代 ListState 或者 ValueState,因为RocksDB 的 map state 并不是将整个 map 作为 value 进行存储,而是将 map 中的一个条目作为键值对进行存储。

如何知道当前 RocksDB 的运行情况

比较直观的方式是打开 RocksDB 的 native metrics ,在默认使用 Flink managed memory 方式的情况下,state.backend.rocksdb.metrics.block-cache-usage ,state.backend.rocksdb.metrics.mem-table-flush-pending,state.backend.rocksdb.metrics.num-running-compactions 以及 state.backend.rocksdb.metrics.num-running-flushes 是比较重要的相关 metrics。

使用 checkpoint 的使用建议

Checkpoint 间隔不要太短

虽然理论上 Flink 支持很短的 checkpoint 间隔,但是在实际生产中,过短的间隔对于底层分布式文件系统而言,会带来很大的压力。另一方面,由于检查点的语义,所以实际上 Flink 作业处理 record 与执行 checkpoint 存在互斥锁,过于频繁的 checkpoint,可能会影响整体的性能。当然,这个建议的出发点是底层分布式文件系统的压力考虑。

合理设置超时时间

默认的超时时间是 10min,如果 state 规模大,则需要合理配置。最坏情况是分布式地创建速度大于单点(job master 端)的删除速度,导致整体存储集群可用空间压力较大。建议当检查点频繁因为超时而失败时,增大超时时间。

【参考文献】

  1. Flink Streaming状态处理(Working with State)

state-management

  • CheckpointedFunction是stateful transformation functions的核心接口,用于跨stream维护state
    • snapshotState 在checkpoint的时候会被调用,用于snapshot state,通常用于flush、commit、synchronize外部系统
    • initializeState 在parallel function初始化的时候(第一次初始化或者从前一次checkpoint recover的时候)被调用,通常用来初始化state,以及处理state recovery的逻辑

从checkpoint中恢复数据时,需要判断snapshot当前的情况,

FunctionSnapshotContext实现了ManagedSnapshotContext, 父类中的方法: getCheckpointId,getCheckpointTimestamp
FunctionInitializationContext实现了ManagedInitializationContext接口, 实现了isRestoredgetOperatorStateStoregetKeyedStateStore方法

在初始化容器之后,我们使用上下文的isrestore()方法检查失败后是否正在恢复。如果是true,即正在恢复,则应用恢复逻辑。

样例: HBase写入OutPutFormat

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
 class PortraitOutputFormat extends RichOutputFormat<EventItem> implements CheckpointedFunction {
// 输出阈值,批量写入的条数
private final int threshold;
// 维护在状态中的数据
private transient ListState<EventItem> checkpointState;
// 内存中的数据
private List<EventItem> bufferedEventItem;
// HBase客户端
private HBaseClient hbaseClient;

public PortraitOutputFormat(HBaseClient hbaseClient) {
this.hbaseClient = hbaseClient;
}
/**
* checkpoint时调用
* 执行snapshot操作,将内存中的数据写入到内存
*/
@Override
public void snapshotState(FunctionSnapshotContext functionSnapshotContext) throws Exception {
checkpointState.clear();
for (EventItem eventItem : bufferedEventItem) {
checkpointState.add(eventItem);
}
}
/**
* 创建state,判断是否存在需要恢复的状态,如果有则需要恢复到bufferedEventItem
*/
@Override
public void initializeState(FunctionInitializationContext context) throws Exception {
ListStateDescriptor<EventItem> descriptor = new ListStateDescriptor<>("buf-p", EventItem.class);
checkpointState = context.getOperatorStateStore().getListState(descriptor);
if (context.isRestored()) {
for (EventItem eventItem : checkpointState.get()) {
bufferedEventItem.add(eventItem);
}
}
}

/**
*
*/
@Override
public void configure(Configuration configuration) {

}
/**
*
*/
@Override
public void open(int taskNumber, int numTasks) throws IOException {
}
/**
* 将新消息写入到缓存bufferedEventItem,缓存个数大约threshold,则执行sink写入,然后清空bufferedEventItem
*/
@Override
public void writeRecord(EventItem value) throws IOException {
if (value.getAttachUserId() == null) {
return;
}

bufferedEventItem.add(value);
int size = bufferedEventItem.size();

if (size >= threshold) {
List<Put> puts = bufferedEventItem
.stream()
.map(eventItem -> {
String rowKey1 = portraitDataGenerator.rowKey(eventItem);
Map<String, String> data = portraitDataGenerator.data(eventItem);
Put put = new Put(rowKey1.getBytes());
for (String cfc : data.keySet()) {
String[] cfcs = cfc.split(":");
String cf = cfcs[0];
String c = cfcs[1];
String dataOne = data.get(cfc);
put.addColumn(cf.getBytes(), c.getBytes(), dataOne.getBytes());
}
return put;
})
.collect(Collectors.toList());
try {
hbaseClient.putAndFlush(puts);
} catch (IOException e) {
e.printStackTrace();
}
bufferedEventItem.clear();
}
}

/**
*
*/
@Override
public void close() throws IOException {
if (hTable != null) {
hTable.flushCommits();
hTable.close();
}
if (connection != null) {
connection.close();
}
}
}

一旦所有checkpoint参与者确认完全,该接口必须由想要接收提交通知的功能/操作来实现。

TTL

1.8 自动清理原理

Apache Flink的1.6.0版本引入了State TTL功能。它使流处理应用程序的开发人员配置过期时间,并在定义时间超时(Time to Live)之后进行清理。在Flink 1.8.0中,该功能得到了扩展,包括对RocksDB和堆状态后端(FSStateBackend和MemoryStateBackend)的历史数据进行持续清理,从而实现旧条目的连续清理过程(根据TTL设置)。

RocksDB后台压缩可以过滤掉过期状态
如果你的Flink应用程序使用RocksDB作为状态后端存储,则可以启用另一个基于Flink特定压缩过滤器的清理策略。RocksDB定期运行异步压缩以合并状态更新并减少存储。Flink压缩过滤器使用TTL检查状态条目的到期时间戳,并丢弃所有过期值。

激活此功能的第一步是通过设置以下Flink配置选项来配置RocksDB状态后端:

state.backend.rocksdb.ttl.compaction.filter.enabled

配置RocksDB状态后端后,将为状态启用压缩清理策略,如以下代码示例所示:

1
2
3
4
StateTtlConfig ttlConfig = StateTtlConfig
.newBuilder(Time.days(7))
.cleanupInRocksdbCompactFilter()
.build();

最佳实践

  • 状态存储的所有数据,均需要考虑清理事迹

【参考文献】

  1. Flink Streaming状态处理(Working with State)

Flink Streaming API

Flink VS Spark

Spark Structured Streaming 是什么?

1、抽象 Abstraction

  Spark中,对于批处理我们有 #RDD ,对于流式,我们有DStream,不过内部实际还是RDD. 所以所有的数据表示本质上还是RDD抽象。
  在Flink中,对于批处理有DataSet,对于流式我们有DataStreams。看起来和Spark类似,他们的不同点在于:

  (一)DataSet在运行时是表现为运行计划(runtime plans)的

  在Spark中,RDD在运行时是表现为java objects的。通过引入Tungsten,这块有了些许的改变。
但是在Flink中是被表现为logical plan(逻辑计划)的, 就是类似于Spark中的dataframes。
所以在Flink中你使用的类Dataframe api是被作为第一优先级来优化的。但是相对来说在Spark RDD中就没有了这块的优化了。
  Flink中的Dataset,对标Spark中的Dataframe,在运行前会经过优化。
在Spark 1.6,dataset API已经被引入Spark了,也许最终会取代RDD 抽象。

  (二)Dataset和DataStream是独立的API

  在Spark中,所有不同的API,例如DStream,Dataframe都是基于RDD抽象的。
  但是在Flink中,Dataset和DataStream是同一个公用的引擎之上两个独立的抽象。所以你不能把这两者的行为合并在一起操作,
  当然,Flink社区目前在朝这个方向努力(https://issues.apache.org/jira/browse/Flink-2320),但是目前还不能轻易断言最后的结果。

2、内存管理

  一直到1.5版本,Spark都是试用java的内存管理来做数据缓存,明显很容易导致OOM或者gc。所以从1.5开始,Spark开始转向精确的控制内存的使用,这就是tungsten项目了。

  而Flink从第一天开始就坚持自己控制内存试用。这个也是启发了Spark走这条路的原因之一。
  Flink除了把数据存在自己管理的内存以外,还直接操作二进制数据。
  在Spark中,从1.5开始,所有的dataframe操作都是直接作用在tungsten的二进制数据上。

3、语言实现

  • 实现语言

    Spark和Flink均有Scala/Java混合编程实现,Spark的核心逻辑由Scala完成,Flink的主要核心逻辑由Java完成

  • 支持应用语言
    Flink主要支持Scala,和Java编程,部分API支持python应用
    Spark主要支持Scala,Java,Python,R语言编程,部分API暂不支持Python和R

4、API

API对比 Flink Spark
应用类型 Batch Streaming Batch Structed Streaming SparkStreaming
数据表示 Dataset datastream RDD,Dataset Dataset Dtream
主要支持API map,filter,flatMap等
转换后数据类型 Dataset datastream RDD,Dataset Dataset Dtream

批处理:

Spark批处理的数据表示经历了从RDD -> DataFrame -> Dataset的变化,均具有不可变,lazy执行,可分区等特性,是Spark框架的核心,rdd经过map等函数操作后,并没有改变而是生成新的RDD,Spark的Dataset(DataFrame是一种特殊的Dataset,已经不推荐使用)还包含数据类型信息

Flink批处理的API是Dataset,同样具有不可变,lazy执行,可分区等特性,是Flink框架的核心,Dataset经过map等函数操作后,并没有改变而是生成新的Dataset

流处理

  • Spark Streaming

    Spark在1.*版本引入的spark streaming作为流处理模块,抽象出Dstream的API来进行流数据处理,同时抽象出通过receiver获取消息数据,然后启动task处理的模式,以及直接启动task消费处理两种方式的流式数据处理。receiver模式由于稳定性不足被遗弃,推荐使用的是直接消费模式;然而本质上讲,Sparkstreaming的流处理是micro-batch的处理模式,将一定时间的流数据作为一个block/RDD,然后使用批处理的rdd的api来完成数据的处理。

  • Structed streaming

    随着Spark在2.*版本的Structed streaming的推出,Spark streaming模块进入了维护模式,从Spark2.*版本以来没有已经没有更新,当前社区主推使用Structed streaming进行流处理。Structed streaming在流处理中有两种流处理模式,一种是microbatch模式;一种是continuous模式;

    • microbatch模式与spark streaming的microbatch模式大致相当,分批处理消息,但可通过设置连续的批次处理,即一个批次执行完之后立即进入下一个批次的处理

    • continuous模式,可以实现真正的流数据处理,端到端的毫秒级,当前处于Experiment状态,也只能支持简单的map,filter操作,当前不支持聚合,current_timestampcurrent_date等操作

    • PS : microbatch <—-> continuous 两种模式可以相互切换且无需改动代码

  • Flink Streaming

    Flink Streaming以流的方式处理流数据,可以实现简单map,fliter等操作,也可以实现复杂的聚合,关联操作,以完善的处理模型及high throughout得到了广泛的应用。

  Spark和Flink都在模仿scala的collection API.所以从表面看起来,两者都很类似。下面是分别用RDD和DataSet API实现的word count

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
// Spark wordcount
object WordCount {

def main(args: Array[String]) {
val env = new SparkContext("local","wordCount")
val data = List("hi","how are you","hi")
val dataSet = env.parallelize(data)
val words = dataSet.flatMap(value => value.split("\\s+"))
val mappedWords = words.map(value => (value,1))
val sum = mappedWords.reduceByKey(_+_)
println(sum.collect())
}
}

// Flink wordcount
object WordCount {

def main(args: Array[String]) {
  val env = ExecutionEnvironment.getExecutionEnvironment
  val data = List("hi","how are you","hi")
  val dataSet = env.fromCollection(data)
  val words = dataSet.flatMap(value => value.split("\\s+"))
  val mappedWords = words.map(value => (value,1))
  val grouped = mappedWords.groupBy(0)
  val sum = grouped.sum(1)
  println(sum.collect())
 }
}

  不知道是偶然还是故意的,API都长得很像,这样很方便开发者从一个引擎切换到另外一个引擎。我感觉以后这种Collection API会成为写data pipeline的标配。

5、Steaming

  Spark把streaming看成是更快的批处理,而Flink把批处理看成streaming的special case。这里面的思路决定了各自的方向,其中两者的差异点有如下这些:

实时 vs 近实时的角度

  Flink提供了基于每个事件的流式处理机制,所以可以被认为是一个真正的流式计算。它非常像storm的model。而Spark,不是基于事件的粒度,而是用小批量来模拟流式,也就是多个事件的集合。所以Spark被认为是近实时的处理系统。

  Spark streaming 是更快的批处理,而Flink Batch是有限数据的流式计算。虽然大部分应用对准实时是可以接受的,但是也还是有很多应用需要event level的流式计算。这些应用更愿意选择storm而非Spark streaming,现在,Flink也许是一个更好的选择。

流式计算和批处理计算的表示

  Spark对于批处理和流式计算,都是用的相同的抽象:RDD,这样很方便这两种计算合并起来表示。而Flink这两者分为了DataSet和DataStream,相比Spark,这个设计算是一个糟糕的设计。

对 windowing 的支持

  因为Spark的小批量机制,Spark对于windowing的支持非常有限。只能基于process time,且只能对batches来做window。而Flink对window的支持非常到位,且Flink对windowing API的支持是相当给力的,允许基于process time,data time,record 来做windowing。我不太确定Spark是否能引入这些API,不过到目前为止,Flink的windowing支持是要比Spark好的。Steaming这部分Flink胜

Window 类型 Window 含义 Flink Streaming SparkStreaming Structed Streaming 备注
tumblingWindow 一个滚动的window 支持 支持 支持
Sliding window 滑动的window 支持 支持 支持
Global window 全局window 支持 间接实现 间接支持 间接支持的含义是可以时间类似功能,但没有抽象出该window
Session window 以接收到数据开始,一定时间没有接收到数据,则结束 支持 不支持 不支持

流join分析:

由于Spark streaming中不支持event time的概念,其只能支持window不同Dstream的RDD的join,不同window间无法join

模块 event-time 流join join实现方式 处理方式 备注
Spark streaming 不支持 支持 window内 processingTime micro-batch处理
FLink1.5之前 支持 支持 window内 native处理,join时(window触发),watermark灵活 Processing Time/ EventTime / element Number
FLink1.6之后 支持 支持 window内,跨window native处理,join时(window触发),watermark灵活 Processing Time/ EventTime / element Number
Structed Streaming 2.2 支持 不支持 仅支持流数据和静态数据的join native处理,join时(window触发),watermark灵活 Processing Time/ EventTime
Structed Streaming 2.3+ 支持 支持 跨window native处理,join时(proocessingTime(interval)触发) Processing Time/ EventTime

PS:

  • Flink/structed streaming开发难度相当,FLink略复杂,但灵活度更高
  • Flink的inteval join
  • Structed Streaming支持数据去重(同个imsi的数据的多个不同join结果的去重)
  • FLink的窗口操作相当于structed streaming的update模式
  • Flink的单流的watermark更新时实时的,有专门线程处理
  • Structed streaming的watermark更新时间基于批的,每个批次共用同一个watermark,如果有多个流,多个流共用一个watermark
  • structed Streaming的watermark更新方法:
    基于每个流找出该流的watermark:Max_event_time - lateness
    找出所有流中最小/最大的watermark设置为batch的watermark
  • Flink专门抽象了类以便不同场景下使用自定义的eventTime的waterMark获取/设置方法,且提供了一般场景下的的类以便使用
  • Flink抽象了trigger和evictor来实现触发计算和清理数据的逻辑,以便自定义相关逻辑
  • FLink 支持sideoutput输出,如迟到的数据可以单独输出

6、SQL interface

  目前Spark-sql是Spark里面最活跃的组件之一,Spark提供了类似Hive的sql和Dataframe这种DSL来查询结构化数据,API很成熟,在流式计算中使用很广,预计在流式计算中也会发展得很快。至于Flink,到目前为止,Flink Table API只支持类似DataFrame这种DSL,并且还是处于beta状态,社区有计划增加SQL 的interface,但是目前还不确定什么时候才能在框架中用上。所以这个部分,Spark胜出。目前Flink已经支持SQL API

7、外部数据源的整合

  Spark的数据源 API是整个框架中最好的,支持的数据源包括NoSql db,parquet,ORC等,并且支持一些高级的操作,例如predicate push down。Flink目前还依赖map/reduce InputFormat来做数据源聚合。这一场Spark胜,目前已经提供

1
2
env.readTextFile(path_i)
env.writeTextFile(path_i)

8、Iterative processing

Flink 迭代处理
Spark迭代处理
  Spark对机器学习的支持较好,因为利用内存cache来加速机器学习算法。然而大部分机器学习算法其实是一个有环的数据流,但是在Spark中,实际是用无环图来表示的,一般的分布式处理引擎都是不鼓励试用有环图的。但是Flink这里又有点不一样,Flink支持在runtime中的有环数据流,这样表示机器学习算法更有效而且更有效率。这一点Flink胜出。

9、Stream as platform vs Batch as Platform

  • Spark诞生在Map/Reduce的时代,数据都是以文件的形式保存在磁盘中,这样非常方便做容错处理。

  • Flink把纯流式数据计算引入大数据时代,无疑给业界带来了一股清新的空气。这个idea非常类似akka-streams这种。

【参考文献】

  1. Apache Flink vs Apache Spark
  2. Flink vs Spark

flink-table-api

官方文档翻译

Concept & Common API

Table API和SQL集成在一个联合的API中。这个API核心概念是Table,
Table可以作为查询的输入和输出。这篇文章展示了使用Table API和SQL查询的通用结构,
如何去进行表的注册,如何去进行表的查询,并且展示如何去进行表的输出。

1. Structure of Table API and SQL Programs

​ 所有使用批量和流式相关的Table API和SQL的程序都有以下相同模式。下面的代码实例展示了Table API和SQL程序的通用结构。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
// 在批处理程序中使用ExecutionEnvironment代替StreamExecutionEnvironment
val env = StreamExecutionEnvironment.getExecutionEnvironment

// 创建TableEnvironment对象
val tableEnv = TableEnvironment.getTableEnvironment(env)

// 注册表
tableEnv.registerTable("table1", ...) // or
tableEnv.registerTableSource("table2", ...) // or
tableEnv.registerExternalCatalog("extCat", ...)

// 基于Table API的查询创建表
val tapiResult = tableEnv.scan("table1").select(...)
// 从SQL查询创建表
val sqlResult = tableEnv.sqlQuery("SELECT ... FROM table2 ...")

// 将表操作API查询到的结果表输出到TableSink,SQL查询到的结果一样如此
tapiResult.writeToSink(...)

// 执行
env.execute()

注意:Table API和SQL查询很容易集成并被嵌入到DataStream或者DataSet程序中。查看将DataStream和DataSet API进行整合章节
学习DataSteams和DataSets是如何转换成Table以及Table是如何转换为DataStream或DataSet

2. Create a TableEnvironment

TableEnvironment是Table API与SQL整合的核心概念之一,它主要有如下功能:

  • 在internal catalog注册表
  • 注册external catalog
  • 执行SQL查询
  • 注册UDF函数(user-defined function),例如 标量, 表或聚合
  • 将DataStream或者DataSet转换为表
  • 保持ExecutionEnvironment或者StreamExecutionEnvironment的引用指向

一个表总是与一个特定的TableEnvironment绑定在一块,
相同的查询不同的TableEnvironment是无法通过join、union合并在一起。

创建TableEnvironment的方法通常是通过StreamExecutionEnvironment,ExecutionEnvironment对象调用其中的静态方法TableEnvironment.getTableEnvironment(),或者是TableConfig来创建。
TableConfig可以用作配置TableEnvironment或是对自定义查询优化器或者是编译过程进行优化(详情查看查询优化)

1
2
3
4
5
6
7
8
9
10
11
12
13
// ***************
// 流式查询
// ***************
val sEnv = StreamExecutionEnvironment.getExecutionEnvironment
// 为流式查询创建一个TableEnvironment对象
val sTableEnv = TableEnvironment.getTableEnvironment(sEnv)

// ***********
// 批量查询
// ***********
val bEnv = ExecutionEnvironment.getExecutionEnvironment
// 为批量查询创建一个TableEnvironment对象
val bTableEnv = TableEnvironment.getTableEnvironment(bEnv)

Register Tables in the Catalog

TableEnvironment包含了通过名称注册表时的表的catalog信息。通常情况下有两种表,一种为输入表,
一种为输出表。输入表主要是在使用Table API和SQL查询时提供输入数据,输出表主要是将Table API和
SQL查询的结果作为输出结果对接到外部系统。

输入表有多种不同的输入源进行注册:

  • 已经存在的Table对象,通常是是作为Table API和SQL查询的结果
  • TableSource,可以访问外部数据如文件,数据库或者是消息系统
  • 来自DataStream或是DataSet程序中的DataStream或DataSet,讨论DataStream或是DataSet
    可以整合DataStream和DataSet API了解到

输出表可使用TableSink进行注册

Register a Table

Table是如何注册到TableEnvironment中如下所示:

1
2
3
4
5
6
7
8
// 获取(创建)TableEnvironment对象
val tableEnv = TableEnvironment.getTableEnvironment(env)

// 从简单的查询结果中作为表
val projTable: Table = tableEnv.scan("X").select(...)

// 将创建的表projTable命名为projectedTable注册到TableEnvironment中
tableEnv.registerTable("projectedTable", projTable)

注意:一张注册过的Table就跟关系型数据库中的视图性质相同,定义表的查询未进行优化,但在另一个查询引用已注册的表时将进行内联。
如果多表查询引用了相同的Table,它就会将每一个引用进行内联并且多次执行,已注册的Table的结果之间不会进行共享。

Register a TableSource

TableSource可以访问外部系统存储例如数据库(Mysql,HBase),特殊格式编码的文件(CSV, Apache [Parquet, Avro, ORC], …)
或者是消息系统 (Apache Kafka, RabbitMQ, …)中的数据。

Flink旨在为通用数据格式和存储系统提供TableSource。请查看此处
了解支持的TableSource类型与如何去自定义TableSour。

TableSource是如何注册到TableEnvironment中如下所示:

1
2
3
4
5
6
7
8
// 获取(创建)TableEnvironment对象
val tableEnv = TableEnvironment.getTableEnvironment(env)

// 创建TableSource对象
val csvSource: TableSource = new CsvTableSource("/path/to/file", ...)

// 将创建的TableSource作为表并命名为csvTable注册到TableEnvironment中
tableEnv.registerTableSource("CsvTable", csvSource)

Register a TableSink

注册过的TableSink可以将SQL查询的结果以表的形式输出到外部的存储系统,例如关系型数据库,
Key-Value数据库(Nosql),消息队列,或者是其他文件系统(使用不同的编码, 例如CSV, Apache [Parquet, Avro, ORC], …)

Flink使用TableSink的目的是为了将常用的数据进行清洗转换然后存储到不同的存储介质中。详情请查看此处
去深入了解哪些sinks是可用的,并且如何去自定义TableSink。

1
2
3
4
5
6
7
8
9
10
11
12
// 获取(创建)TableEnvironment对象
val tableEnv = TableEnvironment.getTableEnvironment(env)

// 创建TableSink对象
val csvSink: TableSink = new CsvTableSink("/path/to/file", ...)

// 定义字段的名称和类型
val fieldNames: Array[String] = Array("a", "b", "c")
val fieldTypes: Array[TypeInformation[_]] = Array(Types.INT, Types.STRING, Types.LONG)

// 将创建的TableSink作为表并命名为CsvSinkTable注册到TableEnvironment中
tableEnv.registerTableSink("CsvSinkTable", fieldNames, fieldTypes, csvSink)

Register an External Catalog

外部目录可以提供有关外部数据库和表的信息,
例如其名称,模式,统计以及有关如何访问存储在外部数据库,表或文件中的数据的信息。

外部目录的创建方式可以通过实现ExternalCatalog接口,并且注册到TableEnvironment中,详情如下所示:

1
2
3
4
5
6
7
8
// 获取(创建)TableEnvironment对象
val tableEnv = TableEnvironment.getTableEnvironment(env)

// 创建一个External Catalog目录对象
val catalog: ExternalCatalog = new InMemoryExternalCatalog

// 将ExternalCatalog注册到TableEnvironment中
tableEnv.registerExternalCatalog("InMemCatalog", catalog)

一旦将External Catalog注册到TableEnvironment中,所有在ExternalCatalog中
定义的表可以通过完整的路径如catalog.database.table进行Table API和SQL的查询操作

目前,Flink提供InMemoryExternalCatalog对象用来做demo和测试,然而,
ExternalCatalog对象还可用作Table API来连接catalogs,例如HCatalog 或 Metastore

Query a Table

Table API

Table API是Scala和Java语言集成查询的API,与SQL查询不同之处在于,它的查询不是像
SQL一样使用字符串进行查询,而是在语言中使用语法进行逐步组合使用

Table API是基于展示表(流或批处理)的Table类,它提供一些列操作应用相关的操作。
这些方法返回一个新的Table对象,该对象表示在输入表上关系运算的结果。一些关系运算是
由多个方法组合而成的,例如 table.groupBy(…).select(),其中groupBy()指定
表的分组,select()表示在分组的结果上进行查询。

Table API
描述了所有支持表的流式或者批处理相关的操作。

下面给出一个简单的实例去说明如何去使用Table API进行聚合查询:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
// 获取(创建)TableEnvironment对象
val tableEnv = TableEnvironment.getTableEnvironment(env)

// 注册Orders表

// 扫描注册过的Orders表
val orders = tableEnv.scan("Orders")

// 计算表中所有来自法国的客户的收入
val revenue = orders
.filter('cCountry === "FRANCE")
.groupBy('cID, 'cName)
.select('cID, 'cName, 'revenue.sum AS 'revSum)

// 将结果输出成一张表或者是转换表

// 执行查询

注意:Scala的Table API使用Scala符号,它使用单引号加字段(‘cID)来表示表的属性的引用,
如果使用Scala的隐式转换的话,确保引入了org.apache.flink.api.scala._ 和 org.apache.flink.table.api.scala._
来确保它们之间的转换。

SQL

Flink的SQL操作基于实现了SQL标准的Apache Calcite,SQL查询通常是使用特殊且有规律的字符串。
SQL
描述了所有支持表的流式或者批处理相关的SQL操作。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// 获取(创建)TableEnvironment对象
val tableEnv = TableEnvironment.getTableEnvironment(env)

// 注册Orders表

// 计算表中所有来自法国的客户的收入
val revenue = tableEnv.sqlQuery("""
|SELECT cID, cName, SUM(revenue) AS revSum
|FROM Orders
|WHERE cCountry = 'FRANCE'
|GROUP BY cID, cName
""".stripMargin)

// 将结果输出成一张表或者是转换表

// 执行查询

下面的例子展示了如何去使用更新查询去插入数据到已注册的表中

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// 获取(创建)TableEnvironment对象
val tableEnv = TableEnvironment.getTableEnvironment(env)

// 注册"Orders"表
// 注册"RevenueFrance"输出表

// 计算表中所有来自法国的客户的收入并且将结果作为结果输出到"RevenueFrance"中
tableEnv.sqlUpdate("""
|INSERT INTO RevenueFrance
|SELECT cID, cName, SUM(revenue) AS revSum
|FROM Orders
|WHERE cCountry = 'FRANCE'
|GROUP BY cID, cName
""".stripMargin)

// 执行查询

Mixing Table API and SQL

Table API和SQL可以很轻松的混合使用因为他们两者返回的结果都为Table对象:

  • 可以在SQL查询返回的Table对象上定义Table API查询
  • 通过在TableEnvironment中注册结果表并在SQL查询的FROM子句中引用它,
    可以在Table API查询的结果上定义SQL查询。

Emit a Table

通过将Table写入到TableSink来作为一张表的输出,TableSink是做为多种文件类型 (CSV, Apache Parquet, Apache Avro),
存储系统(JDBC, Apache HBase, Apache Cassandra, Elasticsearch), 或者是消息系统 (Apache Kafka, RabbitMQ).输出的通用接口,

Batch Table只能通过BatchTableSink来进行数据写入,而Streaming Table可以
选择AppendStreamTableSink,RetractStreamTableSink,UpsertStreamTableSink
中的任意一个来进行。

请查看Table Source & Sinks
来更详细的了解支持的Sinks并且如何去实现自定义的TableSink。

可以使用两种方式来输出一张表:

  • Table.writeToSink(TableSink sink)方法使用提供的TableSink自动配置的表的schema来
    进行表的输出
  • Table.insertInto(String sinkTable)方法查找在TableEnvironment目录中提供的名称下使用特定模式注册的TableSink。
    将输出表的模式将根据已注册的TableSink的模式进行验证

下面的例子展示了如何去查询结果作为一张表输出

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
// 获取(创建)TableEnvironment对象
val tableEnv = TableEnvironment.getTableEnvironment(env)

// 使用Table API或者SQL 查询来查找结果
val result: Table = ...
// 创建TableSink对象
val sink: TableSink = new CsvTableSink("/path/to/file", fieldDelim = "|")

// 方法1: 使用TableSink的writeToSink()方法来将结果输出为一张表
result.writeToSink(sink)

// 方法2: 注册特殊schema的TableSink
val fieldNames: Array[String] = Array("a", "b", "c")
val fieldTypes: Array[TypeInformation] = Array(Types.INT, Types.STRING, Types.LONG)
tableEnv.registerTableSink("CsvSinkTable", fieldNames, fieldTypes, sink)
// 调用注册过的TableSink中insertInto() 方法来将结果输出为一张表
result.insertInto("CsvSinkTable")

// 执行

Translate and Execute a Query

Table API和SQL查询的结果转换为DataStream
或是DataSet
取决于它的输入是流式输入还是批处理输入。查询逻辑在内部表示为逻辑执行计划,并分为两个阶段进行转换:

  • 优化逻辑执行计划
  • 转换为DataStream或DataSet

Table API或SQL查询在下面请看下进行转换:

  • 当调用Table.writeToSink() 或 Table.insertInto()进行查询结果表输出的时候
  • 当调用TableEnvironment.sqlUpdate()进行SQL更新查询时
  • 当表转换为DataSteam或DataSet时,详情查看Integration with DataStream and DataSet API

一旦进行转换后,Table API或SQL查询的结果就会在StreamExecutionEnvironment.execute() 或 ExecutionEnvironment.execute()
被调用时被当做DataStream或DataSet一样被进行处理

Integration with DataStream and DataSet API

Table API或SQL查询的结果很容易被DataStream
或是DataSet内嵌整合。举个例子,
我们会进行外部表的查询(像关系型数据库),然后做像过滤,映射,聚合或者是元数据关联的一些预处理。
然后使用DataStream或是DataSet API(或者是基于这些基础库开发的上层API库, 例如CEP或Gelly)进一步对数据进行处理。
同样,Table API或SQL查询也可以应用于DataStream或DataSet程序的结果。

##implicit Conversion for Scala
Scala Table API具有DataSet,DataStream和Table Class之间的隐式转换,流式操作API中只要引入org.apache.flink.table.api.scala._
和 org.apache.flink.api.scala._ 便可以进行相应的隐式转换

Register a DataStream or DataSet as Table

DataStream或DataSet也可以作为Table注册到TableEnvironment中。结果表的模式取决于已注册的DataStream或DataSet的数据类型,
详情请查看mapping of data types to table schema

1
2
3
4
5
6
7
8
9
10
11
12
// 获取(创建)TableEnvironment对象
// 注册如表一样的DataSet

val tableEnv = TableEnvironment.getTableEnvironment(env)

val stream: DataStream[(Long, String)] = ...

// 将DataStream作为具有"f0", "f1"字段的"myTable"表注册到TableEnvironment中
tableEnv.registerDataStream("myTable", stream)

// 将DataStream作为具有"myLong", "myString"字段的"myTable2"表注册到TableEnvironment中
tableEnv.registerDataStream("myTable2", stream, 'myLong, 'myString)

注意:DataStream表的名称必须与^ DataStreamTable [0-9] +模式不匹配,
并且DataSet表的名称必须与^ DataSetTable [0-9] +模式不匹配。
这些模式仅供内部使用。

Convert a DataStream or DataSet into a Table

如果你使用Table API或是SQL查询,你可以直接将DataStream或DataSet直接转换为表而不需要
再将它们注册到TableEnvironment中。

1
2
3
4
5
6
7
8
9
10
11
12
// 获取(创建)TableEnvironment对象
// 注册如表一样的DataSet
val tableEnv = TableEnvironment.getTableEnvironment(env)

val stream: DataStream[(Long, String)] = ...

// 使用默认的字段'_1, '_2将DataStram转换为Table
val table1: Table = tableEnv.fromDataStream(stream)

// 使用默认的字段'myLong, 'myString将DataStram转换为Table

val table2: Table = tableEnv.fromDataStream(stream, 'myLong, 'myString)

Convert a Table into a DataStream or DataSet

表可以转换为DataStream或DataSet,通过这种方式,自定义DataStream或DataSet
同样也可以作为Table API或SQL查询结果的结果。
当把表转换为DataStream或DataSet时,你需要指定生成的DataStream或DataSet的数据类型。
例如,表格行所需转换的数据类型,通常最方便的转换类型也最常用的是Row。
以下列表概述了不同选项的功能:

  • Row:字段按位置,任意数量的字段映射,支持空值,无类型安全访问。
  • POJO:字段按名称(POJO字段必须与Table字段保持一致),任意数量的字段映射,支持空值,类型安全访问。
  • Case Class:字段按位置,任意数量的字段映射,不支持空值,类型安全访问。
  • Tuple:字段按位置,Scala支持22个字段,Java 25个字段映射,不支持空值,类型安全访问。
  • Atomic Type:表必须具有单个字段,不支持空值,类型安全访问。

Convert a Table into a DataStream

作为流式查询结果的表将动态更新,它随着新记录到达查询的输入流而改变,于是,转换到这样的动态查询DataStream
需要对表的更新进行编码。
将表转换为DataStream有两种模式:

  • Append Mode:这种模式仅用于动态表仅仅通过INSERT来进行表的更新,它是仅可追加模式,
    并且之前输出的表不会进行更改
  • Retract Mode:这种模式经常用到。它使用布尔值的变量来对INSERT和DELETE对表的更新做标记
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
// 获取(创建)TableEnvironment对象 
// 注册如表一样的DataSet
val tableEnv = TableEnvironment.getTableEnvironment(env)

// 表中有两个字段(String name, Integet age)
val table: Table = ...

// 将表转换为列的 append DataStream
val dsRow: DataStream[Row] = tableEnv.toAppendStream[Row](table)

// 将表转换为Tubple2[String,Int]的 append DataStream
// convert the Table into an append DataStream of Tuple2[String, Int]
val dsTuple: DataStream[(String, Int)] dsTuple =
tableEnv.toAppendStream[(String, Int)](table)

// convert the Table into a retract DataStream of Row.
// Retract Mode下将表转换为列的 append DataStream
// 判断A retract stream X是否为DataStream[(Boolean, X)]
// 布尔只表示数据类型的变化,True代表为INSERT,false表示为删除
val retractStream: DataStream[(Boolean, Row)] = tableEnv.toRetractStream[Row](table)

注意:关于动态表和它的属性详情参考Streaming Queries

Convert a Table into a DataSet

表转换为DataSet如下所示:

1
2
3
4
5
6
7
8
9
10
11
12
// 获取(创建)TableEnvironment对象 
// 注册如表一样的DataSet
val tableEnv = TableEnvironment.getTableEnvironment(env)

// 表中有两个字段(String name, Integet age)
val table: Table = ...

// 将表转换为列的DataSet
val dsRow: DataSet[Row] = tableEnv.toDataSet[Row](table)

// 将表转换为Tubple2[String,Int]的DataSet
val dsTuple: DataSet[(String, Int)] = tableEnv.toDataSet[(String, Int)](table)

Mapping of Data Types to Table Schema

Flink的DataStream和DataSet API支持多种类型。组合类型像Tuple(内置Scala元组和Flink Java元组),
POJOs,Scala case classes和Flink中具有可在表表达式中访问的多个字段允许嵌套数据结构的Row类型,
其他类型都被视为原子类型。接下来,我们将会描述Table API是如何将这些类型转换为内部的列展现并且
举例说明如何将DataStream转换为Table

Position-based Mapping

基于位置的映射通常在保持顺序的情况下给字段一个更有意义的名称,这种映射可用于有固定顺序的组合数据类型,
也可用于原子类型。复合数据类型(如元组,行和Case Class)具有此类字段顺序.然而,POJO的字段必须与映射的
表的字段名相同。

当定义基于位置的映射,输入的数据类型不得存在指定的名称,不然API会认为这些映射应该按名称来进行映射。
如果未指定字段名称,则使用复合类型的默认字段名称和字段顺序,或者使用f0作为原子类型。

1
2
3
4
5
6
7
8
9
10
// 获取(创建)TableEnvironment对象 
val tableEnv = TableEnvironment.getTableEnvironment(env)

val stream: DataStream[(Long, Int)] = ...

// 使用默认的字段'_1, '_2将DataStram转换为Table
val table1: Table = tableEnv.fromDataStream(stream)

// 使用默认的字段'myLong, 'myInt将DataStram转换为Table
val table: Table = tableEnv.fromDataStream(stream, 'myLong 'myInt)

Name-based Mapping

基于名称的映射可用于一切数据类型包括POJOs,它是定义表模式映射最灵活的一种方式。虽然查询结果的字段可能会使用别名,但
这种模式下所有的字段都是使用名称进行映射的。使用别名的情况下会进行重排序。
如果未指定字段名称,则使用复合类型的默认字段名称和字段顺序,或者使用f0作为原子类型。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// 获取(创建)TableEnvironment对象 
val tableEnv = TableEnvironment.getTableEnvironment(env)

val stream: DataStream[(Long, Int)] = ...

// 使用默认的字段'_1 和 '_2将DataStram转换为Table
val table: Table = tableEnv.fromDataStream(stream)

// 只使用'_2字段将DataStream转换为Table
val table: Table = tableEnv.fromDataStream(stream, '_2)

// 交换字段将DataStream转换为Table
val table: Table = tableEnv.fromDataStream(stream, '_2, '_1)

// 交换后的字段给予别名'myInt, 'myLong将DataStream转换为Table
val table: Table = tableEnv.fromDataStream(stream, '_2 as 'myInt, '_1 as 'myLong)

Atomic Types

Flink将基础类型(Integer, Double, String)和通用类型(不能被分析和拆分的类型)视为原子类型。
原子类型的DataStream或DataSet转换为只有单个属性的表。从原子类型推断属性的类型,并且可以指定属性的名称。

1
2
3
4
5
6
7
8
9
10
// 获取(创建)TableEnvironment对象
val tableEnv = TableEnvironment.getTableEnvironment(env)

val stream: DataStream[Long] = ...

// 将DataStream转换为带默认字段"f0"的表
val table: Table = tableEnv.fromDataStream(stream)

// 将DataStream转换为带字段"myLong"的表
val table: Table = tableEnv.fromDataStream(stream, 'myLong)

Tuples (Scala and Java) and Case Classes (Scala only)

Flink支持内建的Tuples并且提供了自己的Tuple类给Java进行使用。DataStreams和DataSet这两种
Tuple都可以转换为表。提供所有字段的名称(基于位置的映射)字段可以被重命名。如果没有指定字段的名称,
就使用默认的字段名称。如果原始字段名(f0, f1, … for Flink Tuples and _1, _2, … for Scala Tuples)被引用了的话,
API就会使用基于名称的映射来代替位置的映射。基于名称的映射可以起别名并且会进行重排序。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
// 获取(创建)TableEnvironment对象 
val tableEnv = TableEnvironment.getTableEnvironment(env)

val stream: DataStream[(Long, String)] = ...

// 将默认的字段重命名为'_1,'_2的DataStream转换为Table
val table: Table = tableEnv.fromDataStream(stream)

// 将字段名为'myLong,'myString的DataStream转换为Table(基于位置)
val table: Table = tableEnv.fromDataStream(stream, 'myLong, 'myString)

// 将重排序后字段为'_2,'_1 的DataStream转换为Table(基于名称)
val table: Table = tableEnv.fromDataStream(stream, '_2, '_1)

// 将映射字段'_2的DataStream转换为Table(基于名称)
val table: Table = tableEnv.fromDataStream(stream, '_2)

// 将重排序后字段为'_2给出别名'myString,'_1给出别名'myLong 的DataStream转换为Table(基于名称)
val table: Table = tableEnv.fromDataStream(stream, '_2 as 'myString, '_1 as 'myLong)

// 定义 case class
case class Person(name: String, age: Int)
val streamCC: DataStream[Person] = ...

// 将默认字段'name, 'age的DataStream转换为Table
val table = tableEnv.fromDataStream(streamCC)

// 将字段名为'myName,'myAge的DataStream转换为Table(基于位置)
val table = tableEnv.fromDataStream(streamCC, 'myName, 'myAge)

将重排序后字段为'_age给出别名'myAge,'_name给出别名'myName 的DataStream转换为Table(基于名称)
val table: Table = tableEnv.fromDataStream(stream, 'age as 'myAge, 'name as 'myName)

POJO (Java and Scala)

Flink支持POJO作为符合类型。决定POJO规则的文档请参考这里

当将一个POJO类型的DataStream或者DataSet转换为Table而不指定字段名称时,Table的字段名称将采用JOPO原生的字段名称作为字段名称。
重命名原始的POJO字段需要关键字AS,因为POJO没有固定的顺序,名称映射需要原始名称并且不能通过位置来完成。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
// 获取(创建)TableEnvironment对象
val tableEnv = TableEnvironment.getTableEnvironment(env)

// Person 是一个有两个字段"name"和"age"的POJO
val stream: DataStream[Person] = ...

// 将 DataStream 转换为带字段 "age", "name" 的Table(字段通过名称进行排序)
val table: Table = tableEnv.fromDataStream(stream)

// 将DataStream转换为重命名为"myAge", "myName"的Table(基于名称)
val table: Table = tableEnv.fromDataStream(stream, 'age as 'myAge, 'name as 'myName)

// 将带映射字段'name的DataStream转换为Table(基于名称)
val table: Table = tableEnv.fromDataStream(stream, 'name)

// 将带映射字段'name并重命名为'myName的DataStream转换为Table(基于名称)
val table: Table = tableEnv.fromDataStream(stream, 'name as 'myName)

Row

Row数据类型可以支持任意数量的字段,并且这些字段支持null值。当进行Row DataStream或Row DataSet
转换为Table时可以通过RowTypeInfo来指定字段的名称。Row Type支持基于位置和名称的两种映射方式。
通过提供所有字段的名称可以进行字段的重命名(基于位置),或者是单独选择列来进行映射/重排序/重命名(基于名称)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
// 获取(创建)TableEnvironment对象
val tableEnv = TableEnvironment.getTableEnvironment(env)

// 在`RowTypeInfo`中指定字段"name" 和 "age"的Row类型DataStream
val stream: DataStream[Row] = ...

// 将 DataStream 转换为带默认字段 "age", "name" 的Table
val table: Table = tableEnv.fromDataStream(stream)

// 将 DataStream 转换为重命名字段 'myName, 'myAge 的Table(基于位置)
val table: Table = tableEnv.fromDataStream(stream, 'myName, 'myAge)

// 将 DataStream 转换为重命名字段 'myName, 'myAge 的Table(基于名称)
val table: Table = tableEnv.fromDataStream(stream, 'name as 'myName, 'age as 'myAge)

// 将 DataStream 转换为映射字段 'name的Table(基于名称)
val table: Table = tableEnv.fromDataStream(stream, 'name)

// 将 DataStream 转换为映射字段 'name并重命名为'myName的Table(基于名称)
val table: Table = tableEnv.fromDataStream(stream, 'name as 'myName)

Query Optimization

Apache Flink 基于 Apache Calcite 来做转换和查询优化。当前的查询优化包括投影、过滤下推、
相关子查询和各种相关的查询重写。Flink不去做join优化,但是会让他们去顺序执行(FROM子句中表的顺序或者WHERE子句中连接谓词的顺序)

可以通过提供一个CalciteConfig对象来调整在不同阶段应用的优化规则集,
这个可以通过调用CalciteConfig.createBuilder())获得的builder来创建,
并且可以通过调用tableEnv.getConfig.setCalciteConfig(calciteConfig)来提供给TableEnvironment。

Explaining a Table

Table API为计算Table提供了一个机制来解析逻辑和优化查询计划,这个可以通过TableEnvironment.explain(table)
来完成。它返回描述三个计划的字符串信息:

  • 关联查询抽象语法树,即未优化过的逻辑执行计划
  • 优化过的逻辑执行计划
  • 物理执行计划

下面的实例展示了相应的输出:

1
2
3
4
5
6
7
8
9
10
11
val env = StreamExecutionEnvironment.getExecutionEnvironment
val tEnv = TableEnvironment.getTableEnvironment(env)

val table1 = env.fromElements((1, "hello")).toTable(tEnv, 'count, 'word)
val table2 = env.fromElements((1, "hello")).toTable(tEnv, 'count, 'word)
val table = table1
.where('word.like("F%"))
.unionAll(table2)

val explanation: String = tEnv.explain(table)
println(explanation)

对应的输出如下:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
== 抽象语法树 ==
LogicalUnion(all=[true])
LogicalFilter(condition=[LIKE($1, 'F%')])
LogicalTableScan(table=[[_DataStreamTable_0]])
LogicalTableScan(table=[[_DataStreamTable_1]])

== 优化后的逻辑执行计划 ==
DataStreamUnion(union=[count, word])
DataStreamCalc(select=[count, word], where=[LIKE(word, 'F%')])
DataStreamScan(table=[[_DataStreamTable_0]])
DataStreamScan(table=[[_DataStreamTable_1]])

== 物理执行计划 ==
Stage 1 : Data Source
content : collect elements with CollectionInputFormat

Stage 2 : Data Source
content : collect elements with CollectionInputFormat

Stage 3 : Operator
content : from: (count, word)
ship_strategy : REBALANCE

Stage 4 : Operator
content : where: (LIKE(word, 'F%')), select: (count, word)
ship_strategy : FORWARD

Stage 5 : Operator
content : from: (count, word)
ship_strategy : REBALANCE

Flink用户自定义函数

用户自定义函数是非常重要的一个特征,因为他极大地扩展了查询的表达能力。

在大多数场景下,用户自定义函数在使用之前是必须要注册的。对于Scala的Table API,udf是不需要注册的。
调用TableEnvironment的registerFunction()方法来实现注册。Udf注册成功之后,会被插入TableEnvironment的function catalog,这样table API和sql就能解析他了。
本文会主要讲三种udf:

  • ScalarFunction
  • TableFunction
  • AggregateFunction

1. Scalar Functions 标量函数

标量函数,是指指返回一个值的函数。标量函数是实现讲0,1,或者多个标量值转化为一个新值。

实现一个标量函数需要继承ScalarFunction,并且实现一个或者多个evaluation方法。标量函数的行为就是通过evaluation方法来实现的。evaluation方法必须定义为public,命名为eval。evaluation方法的输入参数类型和返回值类型决定着标量函数的输入参数类型和返回值类型。evaluation方法也可以被重载实现多个eval。同时evaluation方法支持变参数,例如:eval(String… strs)。

下面给出一个标量函数的例子。例子实现的事一个hashcode方法。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
public class HashCode extends ScalarFunction {
private int factor = 12;

public HashCode(int factor) {
this.factor = factor;
}

public int eval(String s) {
return s.hashCode() * factor;
}
}

BatchTableEnvironment tableEnv = TableEnvironment.getTableEnvironment(env);

// register the function
tableEnv.registerFunction("hashCode", new HashCode(10));

// use the function in Java Table API
myTable.select("string, string.hashCode(), hashCode(string)");

// use the function in SQL API
tableEnv.sqlQuery("SELECT string, HASHCODE(string) FROM MyTable");

````

默认情况下evaluation方法的返回值类型是由flink类型抽取工具决定。对于基础类型,简单的POJOS是足够的,但是更复杂的类型,自定义类型,组合类型,会报错。这种情况下,返回值类型的TypeInformation,需要手动指定,方法是重载
ScalarFunction#getResultType()。

下面给一个例子,通过复写ScalarFunction#getResultType(),将long型的返回值在代码生成的时候翻译成Types.TIMESTAMP。

```java
public static class TimestampModifier extends ScalarFunction {
public long eval(long t) {
return t % 1000;
}

public TypeInformation<?> getResultType(signature: Class<?>[]) {
return Types.TIMESTAMP;
}
}

2. Table Functions 表函数

与标量函数相似之处是输入可以0,1,或者多个参数,但是不同之处可以输出任意数目的行数。返回的行也可以包含一个或者多个列。

为了自定义表函数,需要继承TableFunction,实现一个或者多个evaluation方法。表函数的行为定义在这些evaluation方法内部,函数名为eval并且必须是public。TableFunction可以重载多个eval方法。Evaluation方法的输入参数类型,决定着表函数的输入类型。Evaluation方法也支持变参,例如:eval(String… strs)。返回表的类型取决于TableFunction的基本类型。Evaluation方法使用collect(T)发射输出的rows。

在Table API中,表函数在scala语言中使用方法如下:.join(Expression) 或者 .leftOuterJoin(Expression),在java语言中使用方法如下:.join(String) 或者.leftOuterJoin(String)。

Join操作算子会使用表值函数(操作算子右边的表)产生的所有行进行(cross) join 外部表(操作算子左边的表)的每一行。

leftOuterJoin操作算子会使用表值函数(操作算子右边的表)产生的所有行进行(cross) join 外部表(操作算子左边的表)的每一行,并且在表函数返回一个空表的情况下会保留所有的outer rows。

在sql语法中稍微有点区别:
cross join用法是LATERAL TABLE()。
LEFT JOIN用法是在join条件中加入ON TRUE。

下面的理智讲的是如何使用表值函数。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
// The generic type "Tuple2<String, Integer>" determines the schema of the returned table as (String, Integer).
public class Split extends TableFunction<Tuple2<String, Integer>> {
private String separator = " ";

public Split(String separator) {
this.separator = separator;
}

public void eval(String str) {
for (String s : str.split(separator)) {
// use collect(...) to emit a row
collect(new Tuple2<String, Integer>(s, s.length()));
}
}
}

BatchTableEnvironment tableEnv = TableEnvironment.getTableEnvironment(env);
Table myTable = ... // table schema: [a: String]

// Register the function.
tableEnv.registerFunction("split", new Split("#"));

// Use the table function in the Java Table API. "as" specifies the field names of the table.
myTable.join("split(a) as (word, length)").select("a, word, length");
myTable.leftOuterJoin("split(a) as (word, length)").select("a, word, length");

// Use the table function in SQL with LATERAL and TABLE keywords.
join.md
tableEnv.sqlQuery("SELECT a, word, length FROM MyTable, LATERAL TABLE(split(a)) as T(word, length)");
// LEFT JOIN a table function (equivalent to "leftOuterJoin" in Table API).
tableEnv.sqlQuery("SELECT a, word, length FROM MyTable LEFT JOIN LATERAL TABLE(split(a)) as T(word, length) ON TRUE");

需要注意的是PROJO类型不需要一个确定的字段顺序。意味着你不能使用as修改表函数返回的pojo的字段的名字。

默认情况下TableFunction返回值类型是由flink类型抽取工具决定。对于基础类型,简单的POJOS是足够的,但是更复杂的类型,自定义类型,组合类型,会报错。这种情况下,返回值类型的TypeInformation,需要手动指定,方法是重载
TableFunction#getResultType()。

下面的例子,我们通过复写TableFunction#getResultType()方法使得表返回类型是RowTypeInfo(String, Integer)。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
public class CustomTypeSplit extends TableFunction<Row> {
public void eval(String str) {
for (String s : str.split(" ")) {
Row row = new Row(2);
row.setField(0, s);
row.setField(1, s.length);
collect(row);
}
}

@Override
public TypeInformation<Row> getResultType() {
return Types.ROW(Types.STRING(), Types.INT());
}
}

3. Aggregation Functions 聚合函数

用户自定义聚合函数聚合一张表(一行或者多行,一行有一个或者多个属性)为一个标量的值。

上图中是讲的一张饮料的表这个表有是那个字段五行数据,现在要做的事求出所有饮料的最高价。

聚合函数需要继承AggregateFunction。聚合函数工作方式如下:
首先,需要一个accumulator,这个是保存聚合中间结果的数据结构。调用AggregateFunction函数的createAccumulator()方法来创建一个空的accumulator.
随后,每个输入行都会调用accumulate()方法来更新accumulator。一旦所有的行被处理了,getValue()方法就会被调用,计算和返回最终的结果。

对于每个AggregateFunction,下面三个方法都是比不可少的:
createAccumulator()
accumulate()
getValue()

flink的类型抽取机制不能识别复杂的数据类型,比如,数据类型不是基础类型或者简单的pojos类型。所以,类似于ScalarFunction 和TableFunction,AggregateFunction提供了方法去指定返回结果类型的TypeInformation,用的是AggregateFunction#getResultType()。Accumulator类型用的是AggregateFunction#getAccumulatorType()。

除了上面的方法,这里有一些可选的方法。尽管有些方法是让系统更加高效的执行查询,另外的一些在特定的场景下是必须的。例如,merge()方法在会话组窗口上下文中是必须的。当一行数据是被视为跟两个回话窗口相关的时候,两个会话窗口的accumulators需要被join。

AggregateFunction的下面几个方法,根据使用场景的不同需要被实现:
retract():在bounded OVER窗口的聚合方法中是需要实现的。
merge():在很多batch 聚合和会话窗口聚合是必须的。
resetAccumulator(): 在大多数batch聚合是必须的。

AggregateFunction的所有方法都是需要被声明为public,而不是static。定义聚合函数需要实现org.apache.flink.table.functions.AggregateFunction同时需要实现一个或者多个accumulate方法。该方法可以被重载为不同的数据类型,并且支持变参。

在这里就不贴出来AggregateFunction的源码了。

下面举个求加权平均的栗子
为了计算加权平均值,累加器需要存储已累积的所有数据的加权和及计数。在栗子中定义一个WeightedAvgAccum类作为accumulator。尽管,retract(), merge(), 和resetAccumulator()方法在很多聚合类型是不需要的,这里也给出了栗子。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60

/**
* Accumulator for WeightedAvg.
*/
public static class WeightedAvgAccum {
public long sum = 0;
public int count = 0;
}

/**
* Weighted Average user-defined aggregate function.
*/
public static class WeightedAvg extends AggregateFunction<Long, WeightedAvgAccum> {

@Override
public WeightedAvgAccum createAccumulator() {
return new WeightedAvgAccum();
}

@Override
public Long getValue(WeightedAvgAccum acc) {
if (acc.count == 0) {
return null;
} else {
return acc.sum / acc.count;
}
}

public void accumulate(WeightedAvgAccum acc, long iValue, int iWeight) {
acc.sum += iValue * iWeight;
acc.count += iWeight;
}

public void retract(WeightedAvgAccum acc, long iValue, int iWeight) {
acc.sum -= iValue * iWeight;
acc.count -= iWeight;
}

public void merge(WeightedAvgAccum acc, Iterable<WeightedAvgAccum> it) {
Iterator<WeightedAvgAccum> iter = it.iterator();
while (iter.hasNext()) {
WeightedAvgAccum a = iter.next();
acc.count += a.count;
acc.sum += a.sum;
}
}

public void resetAccumulator(WeightedAvgAccum acc) {
acc.count = 0;
acc.sum = 0L;
}
}

// register function
StreamTableEnvironment tEnv = ...
tEnv.registerFunction("wAvg", new WeightedAvg());

// use function
tEnv.sqlQuery("SELECT user, wAvg(points, level) AS avgPoints FROM userScores GROUP BY user");

4. 实现udf的最佳实践经验

Table API和SQL 代码生成器内部会尽可能多的尝试使用原生值。用户定义的函数可能通过对象创建、强制转换(casting)和拆装箱((un)boxing)引入大量开销。因此,强烈推荐参数和返回值的类型定义为原生类型而不是他们包装类型(boxing class)。Types.DATE 和Types.TIME可以用int代替。Types.TIMESTAMP可以用long代替。

我们建议用户自定义函数使用java编写而不是scala编写,因为scala的类型可能会有不被flink类型抽取器兼容。

用Runtime集成UDFs

有时候udf需要获取全局runtime信息或者在进行实际工作之前做一些设置和清除工作。Udf提供了open()和close()方法,可以被复写,功能类似Dataset和DataStream API的RichFunction方法。

Open()方法是在evaluation方法调用前调用一次。Close()是在evaluation方法最后一次调用后调用。

Open()方法提共一个FunctionContext,FunctionContext包含了udf执行环境的上下文,比如,metric group,分布式缓存文件,全局的job参数。

通过调用FunctionContext的相关方法,可以获取到相关的信息:

方法描述

  • getMetricGroup() - 并行子任务的指标组
  • getCachedFile(name) -分布式缓存文件的本地副本
  • getJobParameter(name, defaultValue) - 给定key全局job参数。

下面,给出的例子就是通过FunctionContext在一个标量函数中获取全局job的参数。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
public class HashCode extends ScalarFunction {

private int factor = 0;

@Override
public void open(FunctionContext context) throws Exception {
// access "hashcode_factor" parameter
// "12" would be the default value if parameter does not exist
factor = Integer.valueOf(context.getJobParameter("hashcode_factor", "12"));
}

public int eval(String s) {
return s.hashCode() * factor;
}
}

ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
BatchTableEnvironment tableEnv = TableEnvironment.getTableEnvironment(env);

// set job parameter
Configuration conf = new Configuration();
conf.setString("hashcode_factor", "31");
env.getConfig().setGlobalJobParameters(conf);

// register the function
tableEnv.registerFunction("hashCode", new HashCode());

// use the function in Java Table API
myTable.select("string, string.hashCode(), hashCode(string)");

// use the function in SQL
tableEnv.sqlQuery("SELECT string, HASHCODE(string) FROM MyTable");

内置函数

scala

三元运算符

sql或者table API筛选数据,必须保证每个字段不为空,
Flink内部,中间结果都是通过case class传递,而case class的字段必须保证不能为空

1
2
BOOLEAN.?(VALUE1, VALUE2)
'is_active_user.isNull.?("0", "1")

等值判断

1
'Fuin === 'active_user

scala中的===是运算符重构

  1. Flink Table Api & SQL 翻译目录

本文根据 Apache Flink 系列直播整理而成,由 Apache Flink Contributor、OPPO 大数据平台研发负责人张俊老师分享。主要内容如下:

  • 网络流控的概念与背景

  • TCP的流控机制

  • Flink TCP-based 反压机制(before V1.5)

  • Flink Credit-based 反压机制 (since V1.5)

  • 总结与思考

网络流控的概念与背景

为什么需要网络流控

首先我们可以看下这张最精简的网络流控的图,Producer 的吞吐率是 2MB/s,Consumer 是 1MB/s,这个时候我们就会发现在网络通信的时候我们的 Producer 的速度是比 Consumer 要快的,有 1MB/s 的这样的速度差,假定我们两端都有一个 Buffer,Producer 端有一个发送用的 Send Buffer,Consumer 端有一个接收用的 Receive Buffer,在网络端的吞吐率是 2MB/s,过了 5s 后我们的 Receive Buffer 可能就撑不住了,这时候会面临两种情况:

  • 1.如果 Receive Buffer 是有界的,这时候新到达的数据就只能被丢弃掉了。

  • 2.如果 Receive Buffer 是无界的,Receive Buffer 会持续的扩张,最终会导致 Consumer 的内存耗尽。

网络流控的实现:静态限速

为了解决这个问题,我们就需要网络流控来解决上下游速度差的问题,传统的做法可以在 Producer 端实现一个类似 Rate Limiter 这样的静态限流,Producer 的发送速率是 2MB/s,但是经过限流这一层后,往 Send Buffer 去传数据的时候就会降到 1MB/s 了,这样的话 Producer 端的发送速率跟 Consumer 端的处理速率就可以匹配起来了,就不会导致上述问题。但是这个解决方案有两点限制:

  • 1、事先无法预估 Consumer 到底能承受多大的速率

  • 2、 Consumer 的承受能力通常会动态地波动

网络流控的实现:动态反馈/自动反压

针对静态限速的问题我们就演进到了动态反馈(自动反压)的机制,我们需要 Consumer 能够及时的给 Producer 做一个 feedback,即告知 Producer 能够承受的速率是多少。动态反馈分为两种:

  • 1、负反馈:接受速率小于发送速率时发生,告知 Producer 降低发送速率

  • 2、正反馈:发送速率小于接收速率时发生,告知 Producer 可以把发送速率提上来

让我们来看几个经典案例

案例一:Storm 反压实现

上图就是 Storm 里实现的反压机制,可以看到 Storm 在每一个 Bolt 都会有一个监测反压的线程(Backpressure Thread),这个线程一但检测到 Bolt 里的接收队列(recv queue)出现了严重阻塞就会把这个情况写到 ZooKeeper 里,ZooKeeper 会一直被 Spout 监听,监听到有反压的情况就会停止发送,通过这样的方式匹配上下游的发送接收速率。

案例二:Spark Streaming 反压实现

Spark Streaming 里也有做类似这样的 feedback 机制,上图 Fecher 会实时的从 Buffer、Processing 这样的节点收集一些指标然后通过 Controller 把速度接收的情况再反馈到 Receiver,实现速率的匹配。

疑问:为什么 Flink(before V1.5)里没有用类似的方式实现 feedback 机制?

首先在解决这个疑问之前我们需要先了解一下 Flink 的网络传输是一个什么样的架构。

这张图就体现了 Flink 在做网络传输的时候基本的数据的流向,发送端在发送网络数据前要经历自己内部的一个流程,会有一个自己的 Network Buffer,在底层用 Netty 去做通信,Netty 这一层又有属于自己的 ChannelOutbound Buffer,因为最终是要通过 Socket 做网络请求的发送,所以在 Socket 也有自己的 Send Buffer,同样在接收端也有对应的三级 Buffer。学过计算机网络的时候我们应该了解到,TCP 是自带流量控制的。实际上 Flink (before V1.5)就是通过 TCP 的流控机制来实现 feedback 的。

TCP 流控机制

根据下图我们来简单的回顾一下 TCP 包的格式结构。首先,他有 Sequence number 这样一个机制给每个数据包做一个编号,还有 ACK number 这样一个机制来确保 TCP 的数据传输是可靠的,除此之外还有一个很重要的部分就是 Window Size,接收端在回复消息的时候会通过 Window Size 告诉发送端还可以发送多少数据。

接下来我们来简单看一下这个过程。

TCP 流控:滑动窗口

TCP 的流控就是基于滑动窗口的机制,现在我们有一个 Socket 的发送端和一个 Socket 的接收端,目前我们的发送端的速率是我们接收端的 3 倍,这样会发生什么样的一个情况呢?假定初始的时候我们发送的 window 大小是 3,然后我们接收端的 window 大小是固定的,就是接收端的 Buffer 大小为 5。

首先,发送端会一次性发 3 个 packets,将 1,2,3 发送给接收端,接收端接收到后会将这 3 个 packets 放到 Buffer 里去。

接收端一次消费 1 个 packet,这时候 1 就已经被消费了,然后我们看到接收端的滑动窗口会往前滑动一格,这时候 2,3 还在 Buffer 当中 而 4,5,6 是空出来的,所以接收端会给发送端发送 ACK = 4 ,代表发送端可以从 4 开始发送,同时会将 window 设置为 3 (Buffer 的大小 5 减去已经存下的 2 和 3),发送端接收到回应后也会将他的滑动窗口向前移动到 4,5,6。

这时候发送端将 4,5,6 发送,接收端也能成功的接收到 Buffer 中去。

到这一阶段后,接收端就消费到 2 了,同样他的窗口也会向前滑动一个,这时候他的 Buffer 就只剩一个了,于是向发送端发送 ACK = 7、window = 1。发送端收到之后滑动窗口也向前移,但是这个时候就不能移动 3 格了,虽然发送端的速度允许发 3 个 packets 但是 window 传值已经告知只能接收一个,所以他的滑动窗口就只能往前移一格到 7 ,这样就达到了限流的效果,发送端的发送速度从 3 降到 1。

我们再看一下这种情况,这时候发送端将 7 发送后,接收端接收到,但是由于接收端的消费出现问题,一直没有从 Buffer 中去取,这时候接收端向发送端发送 ACK = 8、window = 0 ,由于这个时候 window = 0,发送端是不能发送任何数据,也就会使发送端的发送速度降为 0。这个时候发送端不发送任何数据了,接收端也不进行任何的反馈了,那么如何知道消费端又开始消费了呢?

TCP 当中有一个 ZeroWindowProbe 的机制,发送端会定期的发送 1 个字节的探测消息,这时候接收端就会把 window 的大小进行反馈。当接收端的消费恢复了之后,接收到探测消息就可以将 window 反馈给发送端端了从而恢复整个流程。TCP 就是通过这样一个滑动窗口的机制实现 feedback。

示例:WindowWordCount

大体的逻辑就是从 Socket 里去接收数据,每 5s 去进行一次 WordCount,将这个代码提交后就进入到了编译阶段。

编译阶段:生成 JobGraph

这时候还没有向集群去提交任务,在 Client 端会将 StreamGraph 生成 JobGraph,JobGraph 就是做为向集群提交的最基本的单元。在生成 JobGrap 的时候会做一些优化,将一些没有 Shuffle 机制的节点进行合并。有了 JobGraph 后就会向集群进行提交,进入运行阶段。

运行阶段:调度 ExecutionGraph

JobGraph 提交到集群后会生成 ExecutionGraph ,这时候就已经具备基本的执行任务的雏形了,把每个任务拆解成了不同的 SubTask,上图 ExecutionGraph 中的 Intermediate Result Partition 就是用于发送数据的模块,最终会将 ExecutionGraph 交给 JobManager 的调度器,将整个 ExecutionGraph 调度起来。然后我们概念化这样一张物理执行图,可以看到每个 Task 在接收数据时都会通过这样一个 InputGate 可以认为是负责接收数据的,再往前有这样一个 ResultPartition 负责发送数据,在 ResultPartition 又会去做分区跟下游的 Task 保持一致,就形成了 ResultSubPartition 和 InputChannel 的对应关系。这就是从逻辑层上来看的网络传输的通道,基于这么一个概念我们可以将反压的问题进行拆解。

问题拆解:反压传播两个阶段

反压的传播实际上是分为两个阶段的,对应着上面的执行图,我们一共涉及 3 个 TaskManager,在每个 TaskManager 里面都有相应的 Task 在执行,还有负责接收数据的 InputGate,发送数据的 ResultPartition,这就是一个最基本的数据传输的通道。在这时候假设最下游的 Task (Sink)出现了问题,处理速度降了下来这时候是如何将这个压力反向传播回去呢?这时候就分为两种情况:

  • 跨 TaskManager ,反压如何从 InputGate 传播到 ResultPartition

  • TaskManager 内,反压如何从 ResultPartition 传播到 InputGate

跨 TaskManager 数据传输

前面提到,发送数据需要 ResultPartition,在每个 ResultPartition 里面会有分区 ResultSubPartition,中间还会有一些关于内存管理的 Buffer。 对于一个 TaskManager 来说会有一个统一的 Network BufferPool 被所有的 Task 共享,在初始化时会从 Off-heap Memory 中申请内存,申请到内存的后续内存管理就是同步 Network BufferPool 来进行的,不需要依赖 JVM GC 的机制去释放。有了 Network BufferPool 之后可以为每一个 ResultSubPartition 创建 Local BufferPool 。 如上图左边的 TaskManager 的 Record Writer 写了 <1,2> 这个两个数据进来,因为 ResultSubPartition 初始化的时候为空,没有 Buffer 用来接收,就会向 Local BufferPool 申请内存,这时 Local BufferPool 也没有足够的内存于是将请求转到 Network BufferPool,最终将申请到的 Buffer 按原链路返还给 ResultSubPartition,<1,2> 这个两个数据就可以被写入了。之后会将 ResultSubPartition 的 Buffer 拷贝到 Netty 的 Buffer 当中最终拷贝到 Socket 的 Buffer 将消息发送出去。然后接收端按照类似的机制去处理将消息消费掉。 接下来我们来模拟上下游处理速度不匹配的场景,发送端的速率为 2,接收端的速率为 1,看一下反压的过程是怎样的。

跨 TaskManager 反压过程

因为速度不匹配就会导致一段时间后 InputChannel 的 Buffer 被用尽,于是他会向 Local BufferPool 申请新的 Buffer ,这时候可以看到 Local BufferPool 中的一个 Buffer 就会被标记为 Used。

发送端还在持续以不匹配的速度发送数据,然后就会导致 InputChannel 向 Local BufferPool 申请 Buffer 的时候发现没有可用的 Buffer 了,这时候就只能向 Network BufferPool 去申请,当然每个 Local BufferPool 都有最大的可用的 Buffer,防止一个 Local BufferPool 把 Network BufferPool 耗尽。这时候看到 Network BufferPool 还是有可用的 Buffer 可以向其申请。

一段时间后,发现 Network BufferPool 没有可用的 Buffer,或是 Local BufferPool 的最大可用 Buffer 到了上限无法向 Network BufferPool 申请,没有办法去读取新的数据,这时 Netty AutoRead 就会被禁掉,Netty 就不会从 Socket 的 Buffer 中读取数据了。

显然,再过不久 Socket 的 Buffer 也被用尽,这时就会将 Window = 0 发送给发送端(前文提到的 TCP 滑动窗口的机制)。这时发送端的 Socket 就会停止发送。

很快发送端的 Socket 的 Buffer 也被用尽,Netty 检测到 Socket 无法写了之后就会停止向 Socket 写数据。

Netty 停止写了之后,所有的数据就会阻塞在 Netty 的 Buffer 当中了,但是 Netty 的 Buffer 是无界的,可以通过 Netty 的水位机制中的 high watermark 控制他的上界。当超过了 high watermark,Netty 就会将其 channel 置为不可写,ResultSubPartition 在写之前都会检测 Netty 是否可写,发现不可写就会停止向 Netty 写数据。

这时候所有的压力都来到了 ResultSubPartition,和接收端一样他会不断的向 Local BufferPool 和 Network BufferPool 申请内存。

Local BufferPool 和 Network BufferPool 都用尽后整个 Operator 就会停止写数据,达到跨 TaskManager 的反压。

TaskManager 内反压过程

了解了跨 TaskManager 反压过程后再来看 TaskManager 内反压过程就更好理解了,下游的 TaskManager 反压导致本 TaskManager 的 ResultSubPartition 无法继续写入数据,于是 Record Writer 的写也被阻塞住了,因为 Operator 需要有输入才能有计算后的输出,输入跟输出都是在同一线程执行, Record Writer 阻塞了,Record Reader 也停止从 InputChannel 读数据,这时上游的 TaskManager 还在不断地发送数据,最终将这个 TaskManager 的 Buffer 耗尽。具体流程可以参考下图,这就是 TaskManager 内的反压过程。

TCP-based 反压的弊端

在介绍 Credit-based 反压机制之前,先分析下 TCP 反压有哪些弊端。

  • 在一个 TaskManager 中可能要执行多个 Task,如果多个 Task 的数据最终都要传输到下游的同一个 TaskManager 就会复用同一个 Socket 进行传输,这个时候如果单个 Task 产生反压,就会导致复用的 Socket 阻塞,其余的 Task 也无法使用传输,checkpoint barrier 也无法发出导致下游执行 checkpoint 的延迟增大。

  • 依赖最底层的 TCP 去做流控,会导致反压传播路径太长,导致生效的延迟比较大。

引入 Credit-based 反压

这个机制简单的理解起来就是在 Flink 层面实现类似 TCP 流控的反压机制来解决上述的弊端,Credit 可以类比为 TCP 的 Window 机制。

Credit-based 反压过程

如图所示在 Flink 层面实现反压机制,就是每一次 ResultSubPartition 向 InputChannel 发送消息的时候都会发送一个 backlog size 告诉下游准备发送多少消息,下游就会去计算有多少的 Buffer 去接收消息,算完之后如果有充足的 Buffer 就会返还给上游一个 Credit 告知他可以发送消息(图上两个 ResultSubPartition 和 InputChannel 之间是虚线是因为最终还是要通过 Netty 和 Socket 去通信),下面我们看一个具体示例。

假设我们上下游的速度不匹配,上游发送速率为 2,下游接收速率为 1,可以看到图上在 ResultSubPartition 中累积了两条消息,10 和 11, backlog 就为 2,这时就会将发送的数据 <8,9> 和 backlog = 2 一同发送给下游。下游收到了之后就会去计算是否有 2 个 Buffer 去接收,可以看到 InputChannel 中已经不足了这时就会从 Local BufferPool 和 Network BufferPool 申请,好在这个时候 Buffer 还是可以申请到的。

过了一段时间后由于上游的发送速率要大于下游的接受速率,下游的 TaskManager 的 Buffer 已经到达了申请上限,这时候下游就会向上游返回 Credit = 0,ResultSubPartition 接收到之后就不会向 Netty 去传输数据,上游 TaskManager 的 Buffer 也很快耗尽,达到反压的效果,这样在 ResultSubPartition 层就能感知到反压,不用通过 Socket 和 Netty 一层层地向上反馈,降低了反压生效的延迟。同时也不会将 Socket 去阻塞,解决了由于一个 Task 反压导致 TaskManager 和 TaskManager 之间的 Socket 阻塞的问题。

总结与思考

总结

  • 网络流控是为了在上下游速度不匹配的情况下,防止下游出现过载

  • 网络流控有静态限速和动态反压两种手段

  • Flink 1.5 之前是基于 TCP 流控 + bounded buffer 实现反压

  • Flink 1.5 之后实现了自己托管的 credit - based 流控机制,在应用层模拟 TCP 的流控机制

思考

有了动态反压,静态限速是不是完全没有作用了?

实际上动态反压不是万能的,我们流计算的结果最终是要输出到一个外部的存储(Storage),外部数据存储到 Sink 端的反压是不一定会触发的,这要取决于外部存储的实现,像 Kafka 这样是实现了限流限速的消息中间件可以通过协议将反压反馈给 Sink 端,但是像 ES 无法将反压进行传播反馈给 Sink 端,这种情况下为了防止外部存储在大的数据量下被打爆,我们就可以通过静态限速的方式在 Source 端去做限流。所以说动态反压并不能完全替代静态限速的,需要根据合适的场景去选择处理方案。

作者:阿里云云栖号
链接:https://juejin.im/post/5dce4b265188254a2b1faddf
来源:掘金
著作权归作者所有。商业转载请联系作者获得授权,非商业转载请注明出处。

Flink Checkpoint与高可用

Flink Checkpoint 受 Chandy-Lamport 分布式快照启发,可以保证数据的高可用。但是有些情况下,不见得一定有效:

Flink On Yarn 模式,某个 Container 发生 OOM 异常,这种情况程序直接变成失败状态,此时 Flink 程序虽然开启 Checkpoint 也无法恢复,因为程序已经变成失败状态,所以此时可以借助外部参与启动程序,比如外部程序检测到实时任务失败时,重新对实时任务进行拉起。

1.1. 2PC

1.1.1. Exactly-once VS At-least-once

算子做快照时,如果等所有输入端的barrier都到了才开始做快照,可保证算子的exactly-once;
如果为了降低延时而跳过对齐,从而继续处理数据,那么等barrier都到齐后做快照就是at-least-once了,因为这次的快照掺杂了下一次快照的数据,当作业失败恢复的时候,这些数据会重复作用系统,就好像这些数据被消费了两遍。

注:对齐只会发生在算子的上端是join操作以及上游存在partition或者shuffle的情况,对于直连操作类似map、flatMap、filter等还是会保证exactly-once的语义。

1.1.2. 端到端的Exactly once实现

2PC分为几个阶段: 开始事务->预提交->提交(或回滚)

为了保证Exactly once, Source和Sink必须支持Flink的2PC

当状态涉及到外部系统时,需要外部系统支持事务操作来配合Flink实现2PC协议,从而保证数据的exatly-once。
这个时候,sink算子除了将自己的state写到状态后端,还必须准备好事务提交。

  • 一旦所有的算子完成了它们的pre-commit,它们会要求一个commit。
  • 如果存在一个算子pre-commit失败了,本次事务失败,我们回滚到上次的checkpoint。
  • 一旦master做出了commit的决定,那么这个commit必须得到执行,就算宕机恢复也有继续执行。

1.1.2.1. pre-commit

pre-commit阶段起始于一次快照的开始,即master节点将checkpoint的barrier注入source端,barrier随着数据向下流动直到sink端。barrier每到一个算子,都会出发算子做本地快照。

precommit

当所有的算子都做完了本地快照并且回复master节点时,pre-commit阶段才算结束。这个时候,checkpoint已经成功,并且包含了外部系统的状态。如果作业失败,可以进行恢复。

precommit-success

1.1.2.2. commit

通知所有的算子这次checkpoint成功了,即2PC的commit阶段。source节点和window节点没有外部状态,所以这时它们不需要做任何操作。
而对于sink节点,需要commit这次事务,将数据写到外部系统。

commit

1.1.2.3. rollback

一旦任何一个算子的快照保存失败,则触发回滚,同样的sink算子也需要取消写入外部的数据

1.1.3. TwoPhaseCommitSinkFunction

为了简化2PC的实现成本,flink抽象了TwoPhaseCommitSinkFunction

  • beginTransaction。开始一次事务,在目的文件系统创建一个临时文件。接下来我们就可以将数据写到这个文件。
  • preCommit。在这个阶段,将文件flush掉,同时重起一个文件写入,作为下一次事务的开始。
  • commit。这个阶段,将文件写到真正的目的目录。值得注意的是,这会增加数据可视的延时。
  • abort。如果回滚,那么删除临时文件。

如果pre-commit成功了,但是commit没有到达算子旧宕机了,flink会将算子恢复到pre-commit时的状态,然后继续commit。

我们需要做的还有就是保证commit的幂等性,这可以通过检查临时文件是否还在来实现。

1.2. checkpoint

保留策略:

  • DELETE_ON_CANCELLATION 表示当程序取消时,删除 Checkpoint 存储文件。
  • RETAIN_ON_CANCELLATION 表示当程序取消时,保存之前的 Checkpoint 存储文件

默认情况下,Flink不会触发一次 Checkpoint 当系统有其他 Checkpoint 在进行时,也就是说 Checkpoint 默认的并发为1。

CheckpointCoordinator :

针对 Flink DataStream 任务,程序需要经历从 StreamGraph -> JobGraph -> ExecutionGraph -> 物理执行图四个步骤,其中在 ExecutionGraph 构建时,会初始化 CheckpointCoordinator。ExecutionGraph通过ExecutionGraphBuilder.buildGraph方法构建,在构建完时,会调用 ExecutionGraph 的enableCheckpointing方法创建CheckpointCoordinator

Flink Checkpoint 参数配置及建议:

  • 当 Checkpoint 时间比设置的 Checkpoint 间隔时间要长时,可以设置 Checkpoint 间最小时间间隔 。这样在上次 Checkpoint 完成时,不会立马进行下一次 Checkpoint,而是会等待一个最小时间间隔,然后在进行该次 Checkpoint。否则,每次 Checkpoint 完成时,就会立马开始下一次 Checkpoint,系统会有很多资源消耗 Checkpoint。
  • 如果Flink状态很大,在进行恢复时,需要从远程存储读取状态恢复,此时可能导致任务恢复很慢,可以设置 Flink Task 本地状态恢复。任务状态本地恢复默认没有开启,可以设置参数state.backend.local-recovery值为true进行激活。
  • Checkpoint保存数,Checkpoint 保存数默认是1,也就是保存最新的 Checkpoint 文件,当进行状态恢复时,如果最新的Checkpoint文件不可用时(比如HDFS文件所有副本都损坏或者其他原因),那么状态恢复就会失败,如果设置 Checkpoint 保存数2,即使最新的Checkpoint恢复失败,那么Flink 会回滚到之前那一次Checkpoint进行恢复。考虑到这种情况,用户可以增加 Checkpoint 保存数。
  • 建议设置的 Checkpoint 的间隔时间最好大于 Checkpoint 的完成时间。

下图是不设置 Checkpoint 最小时间间隔示例图,可以看到,系统一致在进行 Checkpoint,可能对运行的任务产生一定影响:

1.3. savepoint

注意:
使用DataStream进行开发,建议为每个算子定义一个 uid,这样我们在修改作业时,即使导致程序拓扑图改变,由于相关算子 uid 没有变,那么这些算子还能够继续使用之前的状态,如果用户没有定义 uid , Flink 会为每个算子自动生成 uid,如果用户修改了程序,可能导致之前的状态程序不能再进行复用。

Flink 在触发Savepoint 或者 Checkpoint时,会根据这次触发的类型计算出在HDFS上面的目录:

如果类型是 Savepoint,那么 其 HDFS 上面的目录为:Savepoint 根目录+savepoint-jobid前六位+随机数字,具体如下格式:

Checkpoint 目录为 chk-checkpoint ID,具体格式如下:

  • 使用 flink cancel -s 命令取消作业同时触发 Savepoint 时,会有一个问题,可能存在触发 Savepoint 失败。比如实时程序处于异常状态(比如 Checkpoint失败),而此时你停止作业,同时触发 Savepoint,这次 Savepoint 就会失败,这种情况会导致,在实时平台上面看到任务已经停止,但是实际实时作业在 Yarn 还在运行。针对这种情况,需要捕获触发 Savepoint 失败的异常,当抛出异常时,可以直接在 Yarn 上面 Kill 掉该任务。
  • 使用 DataStream 程序开发时,最好为每个算子分配 uid,这样即使作业拓扑图变了,相关算子还是能够从之前的状态进行恢复,默认情况下,Flink 会为每个算子分配 uid,这种情况下,当你改变了程序的某些逻辑时,可能导致算子的 uid 发生改变,那么之前的状态数据,就不能进行复用,程序在启动的时候,就会报错。
  • 由于 Savepoint 是程序的全局状态,对于某些状态很大的实时任务,当我们触发 Savepoint,可能会对运行着的实时任务产生影响,个人建议如果对于状态过大的实时任务,触发 Savepoint 的时间,不要太过频繁。根据状态的大小,适当的设置触发时间。
  • 当我们从 Savepoint 进行恢复时,需要检查这次 Savepoint 目录文件是否可用。可能存在你上次触发 Savepoint 没有成功,导致 HDFS 目录上面 Savepoint 文件不可用或者缺少数据文件等,这种情况下,如果在指定损坏的 Savepoint 的状态目录进行状态恢复,任务会启动不起来。

1.4. snapshot保存到哪里? 应该需要汇总到jobManager?

1.5. state backend

FsStateBackend

构造方法:
FsStateBackend(URI checkpointDataUri,boolean asynchronousSnapshots)

1 基于文件系统的状态管理器
2 如果使用,默认是异步
3 比较稳定,3个副本,比较安全。不会出现任务无法恢复等问题
4 状态大小受磁盘容量限制

存储方式:

  • State: TaskManager内存
  • checkpoint: 外部文件系统(本地或HDFS)

容量限制:

  • 单TaskManager上State总量不超过它的内存
  • 总大小不超过配置的文件系统容量

推荐使用场景:

  • 常规使用状态的作业,例如分钟级窗口聚合、join、窗口比较长、kv状态大;需要开启HA的作业
  • 可以用于生产场景

RocksDBStateBackend

状态数据先写入RocksDB,然后异步的将状态数据写入文件系统。正在进行计算的热数据存储在RocksDB,长时间才更新的数据写入磁盘中(文件系统)存储,体量比较小的元数据状态写入JobManager内存中(将工作state保存在RocksDB中,并且默认将checkpoint数据存在文件系统中)

目前唯一支持incremental的checkpoints的策略

构造方法:
RocksDBStateBackend(URI checkpointDataUri,boolean enableIncrementalCheckpointing)

存储方式:

  • State: TaskManager上的KV数据库(实际使用内存+硬盘)
  • Checkpoint: 外部文件系统(本地或HDFS)

容量限制:

  • 单TaskManager上的State总量不超过他的内存+磁盘
  • 单key最大2G
  • 总大小不超过配置的文件系统容量

推荐使用的场景:

  • 超大状态的作业,例如天级别窗口聚合;需要开启HA的作业;对状态读写性能要求不高的作业
  • 可以在生产环境使用

MemoryStateBackend

构造方法:
MemoryStateBackend(int maxStateSize, boolean asynchronousSnapshots)

主机内存中的数据可能会丢失,任务可能无法恢复

存储方式:

  • State: TaskManager内存
  • Checkpoint: JobManager内存

容量限制

  • 单个State maxStateSize默认5M
  • maxStateSize <= akka.frameSize 默认10M
  • 总大小不超过JobManager的内存

推荐使用场景:

  • 本地测试;几乎无状态的作业,比如ETL;JobManager不容易挂,或挂掉影响不大的情况
  • 不推荐在生产环境使用

1.6. checkpoint 与 savepoint

Checkpoint指定触发生成时间间隔后,每当需要触发Checkpoint时,会向Flink程序运行时的多个分布式的Stream Source中插入一个Barrier标记,这些Barrier会根据Stream中的数据记录一起流向下游的各个Operator。
当一个Operator接收到一个Barrier时,它会暂停处理Steam中新接收到的数据记录。
因为一个Operator可能存在多个输入的Stream,而每个Stream中都会存在对应的Barrier,该Operator要等到所有的输入Stream中的Barrier都到达。(对齐)
当所有Stream中的Barrier都已经到达该Operator,这时所有的Barrier在时间上看来是同一个时刻点(表示已经对齐),在等待所有Barrier到达的过程中,
Operator的Buffer中可能已经缓存了一些比Barrier早到达Operator的数据记录(Outgoing Records),这时该Operator会将数据记录(Outgoing Records)发射(Emit)出去,作为下游Operator的输入,
最后将Barrier对应Snapshot发射(Emit)出去作为此次Checkpoint的结果数据。

Checkpoint 是增量做的,每次的时间较短,数据量较小,只要在程序里面启用后会自动触发,用户无须感知;Checkpoint 是作业 failover 的时候自动使用,不需要用户指定。

Savepoint 是全量做的,每次的时间较长,数据量较大,需要用户主动去触发。Savepoint 一般用于程序的版本更新(详见文档),Bug 修复,A/B Test 等场景,需要用户指定。

保存的内容

  • 首先,Savepoint 包含了一个目录,其中包含(通常很大的)二进制文件,这些文件表示了整个流应用在 Checkpoint/Savepoint 时的状态。
  • 以及一个(相对较小的)元数据文件,包含了指向 Savapoint 各个文件的指针,并存储在所选的分布式文件系统或数据存储中。

目标

Savepoint 和 Checkpoint 的不同之处很像传统数据库中备份与恢复日志之间的区别。Checkpoint 的主要目标是充当 Flink 中的恢复机制,确保能从潜在的故障中恢复。相反,Savepoint 的主要目标是充当手动备份、恢复暂停作业的方法。

实现

Checkpoint 被设计成轻量和快速的机制。它们可能(但不一定必须)利用底层状态后端的不同功能尽可能快速地恢复数据。例如,基于 RocksDB 状态后端的增量检查点,能够加速 RocksDB 的 checkpoint 过程,这使得 checkpoint 机制变得更加轻量。相反,Savepoint 旨在更多地关注数据的可移植性,并支持对作业做任何更改而状态能保持兼容,这使得生成和恢复的成本更高

状态文件保留策略

Checkpoint默认程序删除,可以设置CheckpointConfig中的参数进行保留 。Savepoint会一直保存,除非用户删除 。

应用

  • 部署流应用的一个新版本,包括新功能、BUG 修复、或者一个更好的机器学习模型
  • 引入 A/B 测试,使用相同的源数据测试程序的不同版本,从同一时间点开始测试而不牺牲先前的状态
  • 在需要更多资源时扩容应用程序
  • 迁移流应用程序到 Flink 的新版本上,或者迁移到另一个集群

Flink数据一致性

一、综述

flink 通过内部依赖checkpoint 并且可以通过设置其参数exactly-once 实现其内部的一致性。但要实现其端到端的一致性,还必须保证
1、source 外部数据源可重设数据的读取位置
2、sink端 需要保证数据从故障恢复时,数据不会重复写入外部系统(或者可以逻辑实现写入多次,但只有一次生效的数据sink端)

二、sink 端到端实现方式

幂等操作:
一个操作,可以重复执行多次,但只导致一次结果更改,豁免重复操作执行就不起作用了,他的瑕疵 (在系统恢复的过程中,如果这段时间内多个更新或者插入导致状态不一致,但当数据追上就可以了)
(逻辑与、逻辑或等)具体理解参照自己以前写的文章。
事务写入:
事务应该具有四个属性:原子性、一致性、隔离性、持久性等。其具体的实现方式有两种
(1)、预写日志
简单易于实现,由于数据提前在状态后端中做了缓存,所以无论什么sink系统,都能用这种方式一批搞定,DataStream API提供了一个模板类:GenericWriteAheadSink,来实现这种事务性sink;
缺点:
1)、sink系统没说他支持事务。有可能出现一部分写入集群了。一部分没有写进去(如果实表,再写一次就写重复了)
2)、checkpoint做完了sink才去真正的写入(但其实得等sink都写完checkpoint才能生效,所以WAL这个机制jobmanager确定它写完还不算真正写完,还得有一个外部系统已经确认 完成的checkpoint)
2)、两阶段提交。 flink 真正实现exactle-once
对于每个checkpoint,sink 任务会启动一个事务,并将接下来所有接收的数据添加到事务中,然后将这些数据写入外部sink系统,但不提交他们(这里是预提交)。当checkpoint完成时的通知,它才正式提交事务,实现结果的真正写入;这种方式真正实现了exactly-once,它需要一个提供事务支持的外部sink系统,Flink提供了其具体实现(TwoPhaseCommitSinkFunction接口)

三、 2pc 对外部 sink的要求

1、外部sink系统必须事务支持,或者sink任务必须能够模拟外部系统上的事务;
2、在checkpoint的间隔期间里,必须能够开启一个事务,并接受数据写入。
3、在收到checkpoint完成通知之前,事务必须是“等待提交”的状态,在故障恢复的情况线,这可能需要一些时间。如果个时候sink系统关闭事务(例如超时了),那么未提交的数据就会丢失;
4、四年任务必选能够在进程失败后恢复事务
5、提交事务必须是幂等操作;

四、综上不同Source和sink的一致性保证:

在这里插入图片描述

flink 和kafka 端到端一致性(kafka(source+flink+kafka(sink)))
1、内部 – 利用checkpoint机制,把状态存盘,发生故障的时候可以恢复,保证内部的状态一致性
2、source – kafka consumer作为source,可以将偏移量保存下来,如果后续任务出现了故障,恢复的时候可以由连接器重置偏移量,重新消费数据,保证一致性;

1
2
3
4
kafka 0.8 和kafka 0.11 之后 通过以下配置将偏移量保存,恢复时候重新消费
kafka.setStartFromLatest();
kafka.setCommitOffsetsOnCheckpoints(false);
kafka 0.9 和kafka0.10 未验证是否支持这两个参数(todo)

3、sink FlinkkafkaProducer作为Sink,采用两阶段提交的sink,由下图可以看出flink 0.11 已经默认继承了TwoPhaseCommitSinkFunction
在这里插入图片描述
但我们需要在参数种传入指定语义,它默认时还是at-least-once
此外我们还需要进行一些producer的容错配置:
(1)除了启用Flink的检查点之外,还可以通过将适当的semantic参数传递给FlinkKafkaProducer011(FlinkKafkaProducer对于Kafka> = 1.0.0版本)
(2)来选择三种不同的操作模式
1)、Semantic.NONE 代表at-mostly-once语义
2)、Semantic.AT_LEAST_ONCE(Flink默认设置
3)、Semantic.EXACTLY_ONCE 使用Kafka事务提供一次精确的语义,每当您使用事务写入Kafka时
(3)、请不要忘记消费kafka记录任何应用程序设置所需的设置isolation.leva(read_committed 或者read_uncommitted-后者是默认)
read_committed,只是读取已经提交的数据。

应用;
Semantic.EXACTLY_ONCE依赖与下游系统能支持事务操作.以0.11kafka为例.
transaction.max.timeout.ms 最大超市时长,默认15分钟,如果需要用exactly语义,需要增加这个值。(因为它小于transaction.timeout.ms )
isolation.level 如果需要用到exactly语义,需要在下级consumerConfig中设置read-commited [read-uncommited(默认值)]
transaction.timeout.ms 默认为1hour

其参数对应关系为 和一些报错问题
checkpoint间隔<transaction.timeout.ms<transaction.max.timeout.ms

参考:https://www.cnblogs.com/createweb/p/11971846.html

注意:
1、Semantic.EXACTLY_ONCE 模式每个FlinkKafkaProducer011实例使用一个固定大小的KafkaProducers池。每个检查点使用这些生产者中的每一个。如果并发检查点的数量超过池大小,FlinkKafkaProducer011 将引发异常,并使整个应用程序失败。请相应地配置最大池大小和最大并发检查点数。

2、Semantic.EXACTLY_ONCE采取所有可能的措施,不要留下任何挥之不去的数据,否则这将有碍于消费者更多地阅读Kafka主题。但是,如果flink应用程序在第一个检查点之前失败,则在重新启动此类应用程序后,系统种将没有有关先前池大小信息,因此,在第一个检查点完成前按比例缩小Flink应用程序的FlinkKafkaProducer011.SAFE_SCALE_DOWN_FACTOR

1
2
3
4
5
6
7
8
9
10
//1。设置最大允许的并行checkpoint数,防止超过producer池的个数发生异常
env.getCheckpointConfig.setMaxConcurrentCheckpoints(5)
//2。设置producer的ack传输配置
// 设置超市时长,默认15分钟,建议1个小时以上
producerConfig.put(ProducerConfig.ACKS_CONFIG, 1)
producerConfig.put(ProducerConfig.TRANSACTION_TIMEOUT_CONFIG, 15000)

//3。在下一个kafka consumer的配置文件,或者代码中设置ISOLATION_LEVEL_CONFIG-read-commited
//Note:必须在下一个consumer中指定,当前指定是没用用的
kafkaonfigs.setProperty(ConsumerConfig.ISOLATION_LEVEL_CONFIG,"read_commited")

完整应用代码:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
package com.shufang.flink.connectors

import java.util.Properties
import org.apache.flink.api.common.serialization.SimpleStringSchema
import org.apache.flink.streaming.api.TimeCharacteristic
import org.apache.flink.streaming.api.functions.timestamps.BoundedOutOfOrdernessTimestampExtractor
import org.apache.flink.streaming.api.scala._
import org.apache.flink.streaming.api.windowing.time.Time
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer.Semantic
import org.apache.flink.streaming.connectors.kafka._
import org.apache.flink.streaming.util.serialization.KeyedSerializationSchemaWrapper
import org.apache.kafka.clients.consumer.ConsumerConfig
import org.apache.kafka.clients.producer.ProducerConfig
import org.apache.kafka.common.serialization.StringDeserializer

object KafkaSource01 {
def main(args: Array[String]): Unit = {
val env: StreamExecutionEnvironment = StreamExecutionEnvironment.getExecutionEnvironment

env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime)

//这是checkpoint的超时时间
//env.getCheckpointConfig.setCheckpointTimeout()
//设置最大并行的chekpoint
env.getCheckpointConfig.setMaxConcurrentCheckpoints(5)
env.getCheckpointConfig.setCheckpointInterval(1000) //增加checkpoint的中间时长,保证可靠性


/**
* 为了保证数据的一致性,我们开启Flink的checkpoint一致性检查点机制,保证容错
*/
env.enableCheckpointing(60000)

/**
* 从kafka获取数据,一定要记得添加checkpoint,能保证offset的状态可以重置,从数据源保证数据的一致性
* 保证kafka代理的offset与checkpoint备份中保持状态一致
*/

val kafkaonfigs = new Properties()

//指定kafka的启动集群
kafkaonfigs.setProperty(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092")
//指定消费者组
kafkaonfigs.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "flinkConsumer")
//指定key的反序列化类型
kafkaonfigs.setProperty(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, classOf[StringDeserializer].getName)
//指定value的反序列化类型
kafkaonfigs.setProperty(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, classOf[StringDeserializer].getName)
//指定自动消费offset的起点配置
// kafkaonfigs.setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest")


/**
* 自定义kafkaConsumer,同时可以指定从哪里开始消费
* 开启了Flink的检查点之后,我们还要开启kafka-offset的检查点,通过kafkaConsumer.setCommitOffsetsOnCheckpoints(true)开启,
* 一旦这个检查点开启,那么之前配置的 auto-commit-enable = true的配置就会自动失效
*/
val kafkaConsumer = new FlinkKafkaConsumer[String](
"console-topic",
new SimpleStringSchema(), // 这个schema是将kafka的数据应设成Flink中的String类型
kafkaonfigs
)

// 开启kafka-offset检查点状态保存机制
kafkaConsumer.setCommitOffsetsOnCheckpoints(true)

// kafkaConsumer.setStartFromEarliest()//
// kafkaConsumer.setStartFromTimestamp(1010003794)
// kafkaConsumer.setStartFromLatest()
// kafkaConsumer.setStartFromSpecificOffsets(Map[KafkaTopicPartition,Long]()

// 添加source数据源
val kafkaStream: DataStream[String] = env.addSource(kafkaConsumer)

kafkaStream.print()

val sinkStream: DataStream[String] = kafkaStream.assignTimestampsAndWatermarks(new BoundedOutOfOrdernessTimestampExtractor[String](Time.seconds(5)) {
override def extractTimestamp(element: String): Long = {
element.split(",")(1).toLong
}
})


/**
* 通过FlinkkafkaProduccer API将stream的数据写入到kafka的'sink-topic'中
*/
// val brokerList = "localhost:9092"
val topic = "sink-topic"
val producerConfig = new Properties()
producerConfig.put(ProducerConfig.ACKS_CONFIG, new Integer(1)) // 设置producer的ack传输配置
producerConfig.put(ProducerConfig.TRANSACTION_TIMEOUT_CONFIG, Time.hours(2)) //设置超市时长,默认1小时,建议1个小时以上

/**
* 自定义producer,可以通过不同的构造器创建
*/
val producer: FlinkKafkaProducer[String] = new FlinkKafkaProducer[String](
topic,
new KeyedSerializationSchemaWrapper[String](SimpleStringSchema),
producerConfig,
Semantic.EXACTLY_ONCE
)

// FlinkKafkaProducer.SAFE_SCALE_DOWN_FACTOR
/** *****************************************************************************************************************
* * 出了要开启flink的checkpoint功能,同时还要设置相关配置功能。
* * 因在0.9或者0.10,默认的FlinkKafkaProducer只能保证at-least-once语义,假如需要满足at-least-once语义,我们还需要设置
* * setLogFailuresOnly(boolean) 默认false
* * setFlushOnCheckpoint(boolean) 默认true
* * come from 官网 below:
* * Besides enabling Flink’s checkpointing,you should also configure the setter methods setLogFailuresOnly(boolean)
* * and setFlushOnCheckpoint(boolean) appropriately.
* ******************************************************************************************************************/

producer.setLogFailuresOnly(false) //默认是false


/**
* 除了启用Flink的检查点之外,还可以通过将适当的semantic参数传递给FlinkKafkaProducer011(FlinkKafkaProducer对于Kafka> = 1.0.0版本)
* 来选择三种不同的操作模式:
* Semantic.NONE 代表at-mostly-once语义
* Semantic.AT_LEAST_ONCE(Flink默认设置)
* Semantic.EXACTLY_ONCE:使用Kafka事务提供一次精确的语义,每当您使用事务写入Kafka时,
* 请不要忘记为使用Kafka记录的任何应用程序设置所需的设置isolation.level(read_committed 或read_uncommitted-后者是默认值)
*/

sinkStream.addSink(producer)

env.execute("kafka source & sink")
}
}

[参考文献]

  1. Flink Checkpoint、Savepoint配置与实践
  2. Flink 小贴士 (2):Flink 如何管理 Kafka 消费位点
  3. Flink实时计算-深入理解Checkpoint和Savepoint
  4. Lightweight Asynchronous Snapshots for Distributed Dataflows: 分布式数据流轻量级异步快照

Flink Checkpoint 原理流程以及常见失败原因分析

影响Checkpoint的几个关键参数:

参数 默认值 备注
state.backend none 用于指定checkpoint state存储的backend,
state.backend.async true 用于指定backend是否使用异步snapshot,有些不支持async或者只支持async的state backend可能会忽略这个参数
state.backend.fs.memory-threshold 1024 用于指定存储于files的state大小阈值,如果小于该值则会存储在root checkpoint metadata file
state.backend.incremental false 用于指定是否采用增量checkpoint,有些不支持增量checkpoint的backend会忽略该配置
state.backend.local-recovery false 任务本地恢复
state.checkpoints.dir none 用于指定checkpoint的data files和meta data存储的目录,该目录必须对所有参与的TaskManagers及JobManagers可见
state.checkpoints.num-retained 1 用于指定保留的已完成的checkpoints个数
state.savepoints.dir none 用于指定savepoints的默认目录
taskmanager.state.local.root-dirs none

RETAIN_ON_CANCELLATION

CheckpointConfig config = env.getCheckpointConfig();
config.enableExternalizedCheckpoints(ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);
config.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
config.setCheckpointInterval(``60000``);

ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION,表示一旦Flink处理程序被cancel后,会保留Checkpoint数据,以便根据实际需要恢复到指定的Checkpoint处理

简单之美 | Flink Checkpoint、Savepoint配置与实践