0%

Flink-sink

Sink的三种模式

Flink table的三种sink模式

Sink有INSERT、UPDATE 和 DELETE 三类,Table的sink模式有append、upsert和retract三种

Sink模式 Insert Update Delete 支持的存储
Append(追加) 支持 不支持 不支持
Upsert(重复时更新) 支持 支持 支持 KV: HBase JDBC
Retract(允许撤销) 支持 支持 支持

Upsert模式与Retract模式的区别

Upsert模式需要唯一的key来传递更新消息,外部连接器需要明确知道这个唯一key的属性

Upsert模式和Retract模式

消息: 都为(Boolean,Row)二元组, 第一个元素代表操作类型:

模式 操作类型 插入 更新 删除
Upsert模式 true 为 UPSERT消息(不存在则INSERT, 存在则UPDATE)
false 为 DELETE消息
upsert消息 upsert消息 delete消息
Retract模式 true为添加消息
false为撤回消息
添加消息 已更新行(上一行)的撤回消息
更新行(新行)的添加消息
撤回消息

Append模式-窗口聚合中的应用

在实时聚合统计中,聚合统计的结果输出是由 Trigger 决定的,而 Append-Only 则意味着对于每个窗口实例(Pane,窗格)Trigger 只能触发一次,则就导致无法在迟到数据到达时再刷新结果。

通常来说,我们可以给 Watermark 设置一个较大的延迟容忍阈值来避免这种刷新(再有迟到数据则丢弃),但代价是却会引入较大的延迟。

Upsert模式

支持 Append-Only 的操作和在有主键的前提下的 Update 和 Delete 操作.

重复时更新

Upsert 模式依赖业务主键来实现输出结果的更新和删除,因此非常适合 KV 数据库,比如 HBase、JDBC 的 TableSink 都使用了这种方式。

Upsert 模式是目前来说比较实用的模式,因为大部分业务都会提供原子或复合类型的主键,而在支持 KV 的存储系统也非常多,但要注意的是不要变更主键,具体原因会在下一节谈到。

Retract模式

允许撤销

举个例子,假设我们将电商订单按照承运快递公司进行分类计数,有如下的结果表。

公司 订单数
中通 2
圆通 1
顺丰 3

那么如果原本一单为中通的快递,后续更新为用顺丰发货,对于 Upsert 模式会产生 (true, (顺丰, 4)) 这样一条 changelog,但中通的订单数没有被修正。相比之下,Retract 模式产出 (false, (中通, 1))(true, (顺丰, 1)) 两条数据,则可以正确地更新数据。