Skip to content

FlinkSQL 的 retract 机制:撤回消息为什么由上游算子发

StreamingSQL和BatchSQL的流程对比

以下是一个嵌套的FlinkSQL代码:

sql
SELECT
cnt,
count(cnt) AS freq
FROM(
  SELECT word, COUNT(num) AS cnt
  FROM Table GROUP BY word
) GROUP BY cnt;
  1. 计算不同的单词出现的频率(task 1)
  2. 计算不同频率下的单词的个数(task 2)

原始的数据:

wordnum
Hello1
Bob1
Word1
Hello1

batch的结果输出:

计算不同单词的出现频率, task1的输出如下所示:

wordnum
Hello2
Bob1
World1

计算不同频率下的个数, task2的输出如下所示:

cntfreq
12
21

streaming的结果输出:

对于Streaming作业来说,数据是源源不断的, 因此写到下游的结果也是源源不断的. 因此,对于StreamingSQL,无法像batchSQL那样只产生一次的结果输出.

sourcetask 1task 2
消息编号wordnumwordnumcntfreq
1hello1hello111
2word1word112
3bob1bob113
4hello1hello221
12

如上表格所示为不同的消息进入系统之后,各个task产生的输出情况. 前三条消息到来之后, task1 和task2的输出比较容易理解,是正常的累加逻辑. 重点看第四条消息到来时的各个task的输出: task 1的输出为 hello, 2 , 这里比较好理解, 由于是按字段进行累计和count, task1缓存了 hello, 1的状态, 当 hello, 1消息进入时,进行累计计算, 输出为 hello, 2. task 2 在接受到 (hello,2)之后, 输出 ( 2, 1) , 并同时产生(1, 2) , 用来覆盖掉之前的( 1, 3).

retract机制实现原理解析

撤回消息由上游算子产生

通过上面的结果输出, 我们大致能明白, Streaming作业和batch作业两种作业的差异, Streaming作业的结果会根据当前的实时数据不断的去修正最终的结果.

其中的关键问题就是: task 2 怎样才能输出 (1 , 2) 这条结果 ? 可能会存在两种方案:

  • task2接受到(hello, 2)这条消息之后, 通过内部的状态信息得出需要减去(hello, 1)这条消息, 将(1, 3) 减去 (hello, 1),得到(1, 2).
  • task2需要接受到上游发送过来的-(hello, 1)的消息, 将 (1, 3)减去(hello, 1), 得到 (1, 2)

所以以上的问题就变成了: 是由task1 还是task2来产生 -(hello, 1) 这条消息 ?

  • 由task2来产生减(hello, 1)的消息
  • 由task1来产生减(hello, 1)的消息

task2产生 -(hello, 1)的消息

task2中保存的状态

  1. Map<String, Integer> 保存所有 word以及对应的count
  2. Map<Integer, Integer> 保存单词出现个数以及对应的频率

task1保存的状态

  1. Map<String, Integer> 保存所有word对应的count.

task1产生-(hello, 1)的消息

task2中保存的状态

  1. Map<Integer, Integer> 保存单词出现的个数以及对应的频率

task1中保存的状态

  1. Map<String, Integer> 保存所有的word以及对应的count

很明显, task1中已经保存了所有的word对应的count, 则task2也不需要进行保存, 由task1产生 减(hello, 1)的所需要保存的状态较少.

因此, task1在接收到第四条消息时,需要产生两条消息: 6. - (hello, 1) 7. +(hello, 2)

什么场景下需要retract

简单SQL

sql
SELECT word, num % 10 as num AS cnt
FROM Table;

简单的字段转换,或者映射, 不涉及到task之间,task和外部系统的数据更新操作, 则不需要retract.

聚合SQL

sql
SELECT
cnt,
count(word) AS freq
FROM(
  SELECT word, COUNT(num) AS cnt
  FROM Table GROUP BY word
);

聚合操作(非时间窗口),涉及到task和task之间,task和外部系统之间的数据更新, 则需要retract机制.

Flink框架本身支持处理和产生add, update, delete类型的消息, 同时也需要最终的sink也需要能够处理这些类型的消息, 如果不能支持, 则有些场景下,就可能无法支持.

类型特点相关系统
append只能接受append消息, 无法处理update,delete消息消息队列(Kafka), druid, opentsdb
upsert可以处理 add, update, delete消息mysql, hbase, kv, es, kudu, es
retract只能处理 add, delete消息, 无法处理update消息print(测试用)

tips:

  • 在实际的应用中, Kafka并不仅仅只能作为append表, 虽然Kafka系统本身无法处理delete或者update消息, 但是在实现上, 可以将 append, delete, update等消息的类型也一并写入到消息体中, 由下游再去处理不同类型的消息类型即可, 实现细节可以参考FlinkKafka Retract-Table支持写retract信息到下游Kafka
  • 很少有系统真正是retract表, 一般支持删除的系统都支持处理update消息.
  • retract无法处理update消息, 如果下游是retract表, 那么Flink框架会将update的消息转化为 delete + add 消息.
  • ~~如果使用了Upsert类型的sink表, 一定要使用 insert into SinkTable select xxx, sum(xxx) group by xxx的写法, 让框架能够识别到sink表的主键用于优化生成的DAG图. (~~在Flink1.12中, 声明主键即可)

本文是「我的数据空间」实时计算实践笔记的一篇。平台把这类链路做成了托管作业(Flink SQL 流作业、实时集成、Iceberg 表管理与维护),支持私有化部署与 OEM 合作 —— 见核心能力总览技术博客,交流合作 QQ:1559851993。