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)) 两条数据,则可以正确地更新数据。