ARTICLE DETAIL

资讯详情

深耕商务建站与企业官网运营的一线实战洞察。

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

FlinkSQL 的 retract 机制:撤回消息为什么由上游算子发 「我的数据空间」实时计算实践笔记 · Flink SQL 系列StreamingSQL和BatchSQL的流程对比以下是一个嵌套的FlinkSQL代码:SELECTcnt,count(cnt)ASfreqFROM(SELECTword,COUNT(num)AScntFROMTableGROUPBYword)GROUPBYcnt;计算不同的单词出现的频率(task 1)计算不同频率下的单词的个数(task 2)原始的数据:wordnumHello1Bob1Word1Hello1batch的结果输出:计算不同单词的出现频率, task1的输出如下所示:wordnumHello2Bob1World1计算不同频率下的个数, task2的输出如下所示:cntfreq1221streaming的结果输出:对于Streaming作业来说数据是源源不断的 因此写到下游的结果也是源源不断的. 因此对于StreamingSQL无法像batchSQL那样只产生一次的结果输出.sourcetask 1task 2消息编号wordnumwordnumcntfreq1hello1hello1112word1word1123bob1bob1134hello1hello22112如上表格所示为不同的消息进入系统之后各个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中保存的状态MapString, Integer 保存所有 word以及对应的countMapInteger, Integer 保存单词出现个数以及对应的频率task1保存的状态MapString, Integer 保存所有word对应的count.task1产生-(hello, 1)的消息task2中保存的状态MapInteger, Integer 保存单词出现的个数以及对应的频率task1中保存的状态MapString, Integer 保存所有的word以及对应的count很明显 task1中已经保存了所有的word对应的count 则task2也不需要进行保存 由task1产生 减(hello, 1)的所需要保存的状态较少.因此 task1在接收到第四条消息时需要产生两条消息:6. - (hello, 1)7. (hello, 2)什么场景下需要retract简单SQLSELECTword,num%10asnumAScntFROMTable;简单的字段转换或者映射 不涉及到task之间task和外部系统的数据更新操作 则不需要retract.聚合SQLSELECTcnt,count(word)ASfreqFROM(SELECTword,COUNT(num)AScntFROMTableGROUPBYword);聚合操作(非时间窗口)涉及到task和task之间task和外部系统之间的数据更新 则需要retract机制.Flink框架本身支持处理和产生add, update, delete类型的消息 同时也需要最终的sink也需要能够处理这些类型的消息 如果不能支持 则有些场景下就可能无法支持.类型特点相关系统append只能接受append消息 无法处理updatedelete消息消息队列(Kafka), druid, opentsdbupsert可以处理 add, update, delete消息mysql, hbase, kv, es, kudu, esretract只能处理 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中 声明主键即可本文收录于「我的数据空间」技术库——一套可私有化部署的数据平台(数据集成 / 实时计算 / 数据湖 / 湖仓查询 / 智能问数)。产品介绍见我的数据空间官网,支持私有化部署与 OEM 合作。
返回列表
PREV
查看更多资讯
NEXT
返回资讯列表