
本文主要介绍了Apache Flink的分组聚合、over聚合和windowjoin操作,以及它们的限制和性能调优。Flink支持标准的GROUPBY子句进行数据分组,还提供了ROLLUP、CUBE和HAVING等快捷方式来表示常见类型的分组集。OVER聚合可以计算有序行范围内每个输入行的聚合值,而窗口联接可以将时间维度添加到联接条件本身中。目前(截至版本Flink1.17),窗口联接要求连接条件包含输入表的窗口开始相等和输入表的窗口结束相等。将来,我们还可以简化joinon子句,如果窗口TVF是TUMBLE或HOP,则只包含windowstart相等。
1、Flink部署、概念介绍、source、transformation、sink使用示例、四大基石介绍和示例等系列综合文章链接
13、Flink的tableapi与sql的基本概念、通用api介绍及入门示例
14、Flink的tableapi与sql之数据类型:内置数据类型以及它们的属性
15、Flink的tableapi与sql之流式概念-详解的介绍了动态表、时间属性配置(如何处理更新结果)、时态表、流上的join、流上的确定性以及查询配置
16、Flink的tableapi与sql之连接外部系统:读写外部系统的连接器和格式以及FileSystem示例(1)
16、Flink的tableapi与sql之连接外部系统:读写外部系统的连接器和格式以及Elasticsearch示例(2)
16、Flink的tableapi与sql之连接外部系统:读写外部系统的连接器和格式以及ApacheKafka示例(3)
16、Flink的tableapi与sql之连接外部系统:读写外部系统的连接器和格式以及JDBC示例(4)
16、Flink的tableapi与sql之连接外部系统:读写外部系统的连接器和格式以及ApacheHive示例(6)
20、FlinkSQL之SQLClient:不用编写代码就可以尝试FlinkSQL,可以直接提交SQL任务到集群上
22、Flink的tableapi与sql之创建表的DDL24、Flink的tableapi与sql之Catalogs
26、Flink的SQL之概览与入门示例
27、Flink的SQL之SELECT(select、where、distinct、orderby、limit、集合操作和去重)介绍及详细示例(1)
27、Flink的SQL之SELECT(SQLHints和Joins)介绍及详细示例(2)
27、Flink的SQL之SELECT(窗口函数)介绍及详细示例(3)
27、Flink的SQL之SELECT(窗口聚合)介绍及详细示例(4)
27、Flink的SQL之SELECT(GroupAggregation分组聚合、OverAggregationOver聚合和WindowJoin窗口关联)介绍及详细示例(5)
30、FlinkSQL之SQL客户端(通过kafka和filesystem的例子介绍了配置文件使用-表、视图等)
41、Flink之Hive方言介绍及详细示例
42、Flink的tableapi与sql之HiveCatalog
43、Flink之Hive读写及详细验证示例
44、Flink之module模块介绍及使用示例和FlinkSQL使用hive内置函数及自定义函数详细示例–网上有些说法好像是错误的
本文介绍了Flink的分组聚合、over聚合和windowjoin操作、当前版本的限制以及具体的运行示例。
本文依赖flink和kafka集群能正常使用。
本文分为3个部分,即介绍了Flink分组聚合、over聚合以及windowjoin,并且每个内容均以验证通过示例进行说明。
本文运行环境是Flink1.17版本。
像大多数数据系统一样,ApacheFlink支持聚合函数:内置和用户定义。用户定义的函数必须在使用前在目录中注册。
聚合函数从多个输入行计算单个结果。例如,有一些聚合可以计算一组行的COUNT,SUM,AVG(average),MAX(maximum)和MIN(minimum)。
下文用到的数据源为orders3,其数据结构以及数据为
对于流式查询,重要的是要了解Flink运行永不终止的连续查询。相反,他们根据其输入表上的更新更新其结果表。对于上面的查询,Flink会在每次将新行插入Orders3表时输出更新的计数。
ApacheFlink支持用于聚合数据的标准GROUPBY子句。
对于流式处理查询,计算查询结果所需的状态可能会无限增长。状态大小取决于组的数量以及聚合函数的数量和类型。例如,MIN/MAX在状态大小上很重,而COUNT很便宜。您可以为查询配置提供适当的状态生存时间(TTL),以防止状态大小过大。请注意,这可能会影响查询结果的正确性。关于TTL的配置可以参考文章:43、Flink之Hive读写及详细验证示例
ApacheFlink为群聚合提供了一套性能调优方式,详见45、Flink之性能调优介绍及示例。该篇文章中将详细介绍性能调优的几个方向及示例。
distinct聚合是删除重复值在聚合时。下面的示例计算orders3表中非重复u_id数,而不是行总数。
对于流式处理查询,计算查询结果所需的状态可能会无限增长。状态大小主要取决于不同行的数量和维护组的时间,按窗口划分的短期组不是问题(shortlivedgroupbywindowsarenotaproblem)。您可以为查询配置提供适当的状态生存时间(TTL),以防止状态大小过大。请注意,这可能会影响查询结果的正确性。关于TTL的配置可以参考文章:43、Flink之Hive读写及详细验证示例
分组集允许比标准GROUPBY描述的操作更复杂的分组操作。行按每个指定的分组集单独分组,并为每个组计算聚合,就像简单的GROUPBY子句一样。
关于groupingset的更多内容可以参考文章:27、Flink的SQL之SELECT(窗口聚合)介绍及详细示例(4)
GROUPINGSETS的每个子列表可以指定零个或多个列或表达式,并且解释方式与直接在GROUPBY子句中使用的方式相同。空分组集意味着所有行都聚合到单个组,即使不存在输入行,也会输出该组。
对分组列或表达式的引用将替换为结果行中的null值,用于未显示这些列的分组集。
对于流式处理查询,计算查询结果所需的状态可能会无限增长。状态大小取决于组集的数量和聚合函数的类型。您可以为查询配置提供适当的状态生存时间(TTL),以防止状态大小过大。请注意,这可能会影响查询结果的正确性。关于TTL的配置可以参考文章:43、Flink之Hive读写及详细验证示例
ROLLUP是用于指定常见类型的分组集的速记表示法(shorthandnotation)。它表示给定的表达式列表和列表的所有前缀,包括空列表。
下面两个查询等效。
CUBE是用于指定常见类型的分组集的速记表示法。它表示给定的列表及其所有可能的子集-幂集。
以下两个查询是等效的。
HAVING消除不满足条件的组行。HAVING与WHERE不同:WHERE在GROUPBY之前筛选单个行,而HAVING筛选由GROUPBY创建的组行。条件中引用的每个列都必须明确引用分组列,除非它出现在聚合函数中。
HAVING的存在会将查询转换为分组查询,即使没有GROUPBY子句也是如此。这与查询包含聚合函数但没有GROUPBY子句时发生的情况相同。该查询将所有选定的行视为形成一个组,并且SELECT列表和HAVING子句只能引用聚合函数中的表列。如果HAVING条件为真,则此类查询将发出单行,如果条件不为真,则发出零行。
OVER聚合计算有序行范围内每个输入行的聚合值。与GROUPBY聚合相比,OVER聚合不会将每个组的结果行数减少到一行。相反,OVER聚合为每个输入行生成一个聚合值。
可以在SELECT子句中定义多个OVER窗口聚合。但是,对于流式处理查询,由于当前限制,所有聚合的OVER窗口必须相同。
OVER窗口是在有序的行序列上定义的。由于表没有固有的顺序,因此ORDERBY子句是必需的。对于流式查询,Flink目前(截至Flink版本1.17)仅支持使用升序时间属性顺序定义的OVER窗口。不支持其他排序。
可以在分区表上定义OVER窗口。在存在PARTITIONBY子句的情况下,仅针对其分区的行计算每个输入行的聚合。
rangedefinition指定聚合中包含多少行。该范围由定义下限和上限的BETWEEN子句定义。这些边界之间的所有行都包含在聚合中。Flink仅支持CURRENTROW作为上限。
有两个选项可以定义范围:行间隔和范围间隔。
RANGE间隔是在ORDERBY列的值上定义的,在Flink的情况下,它始终是一个时间属性。以下RANGE间隔定义时间属性最多比当前行少30分钟的所有行都包含在聚合中。
行间隔(ROWSinterval)是基于计数的间隔。它准确定义聚合中包含的行数。以下ROWS间隔定义聚合中包括当前行和当前行之前的10行(因此总共11行)。
WINDOW子句可用于在SELECT子句之外定义OVER窗口。它可以使查询更具可读性,还允许我们为多个聚合重用窗口定义。
以下查询计算每个订单的当前订单前一小时内收到的同一用户的所有订单的金额总和。
其实用的是proctime,仅仅是示例,可能实际的业务上来说不够准确,仅仅是为了模拟验证。
窗口联接(windowjoin)将时间维度添加到联接条件本身中。在此过程中,窗口联接将联接共享公共键且位于同一窗口中的两个流的元素。窗口联接的语义与数据流窗口联接相同。
对于流式处理查询,与连续表上的其他联接不同,窗口联接不会发出中间结果,而只会在窗口末尾发出最终结果。此外,窗口联接在不再需要时清除所有中间状态。
通常,窗口联接与WindowingTVF一起使用。此外,窗口联接可以遵循基于WindowingTVF的其他操作,例如窗口聚合,窗口TopN和窗口连接。
目前(截至版本Flink1.17),窗口联接要求连接条件包含输入表的窗口开始相等和输入表的窗口结束相等。
窗口联接支持INNER/LEFT/RIGHT/FULLOUTER/ANTI/SEMIJOIN。
下面显示了INNER/LEFT/RIGHT/FULLOUTERWindowJoin语句的语法。
INNER/LEFT/RIGHT/FULLOUTERWINDOWJOIN的语法彼此非常相似,我们在这里只举一个FULLOUTERJOIN的例子。执行窗口联接时,具有公共键和公共滚动窗口的所有元素将联接在一起。我们只给出一个适用于滚动窗口TVF的窗口联接示例。通过将联接的时间区域范围界定为固定的五分钟间隔,我们将数据集切成两个不同的时间窗口:[9:35,9:40)和[9:40,9:45)。L3和R3行无法联接在一起,因为它们属于不同的窗口。
SemiWindowJoinsreturnsarowfromoneleftrecordifthereisatleastonematchingrowontherightsidewithinthecommonwindow.如果公共窗口中右侧至少有一个匹配的行,则半窗口联接(SemiWindowJoins)从左侧记录返回一行。
反窗口联接(AntiWindowJoins)是内部窗口联接的正面:它们包含每个公共窗口中所有未联接的行。
目前(截至FLink版本1.17),窗口连接要求连接条件包含输入表的windowstarts相等和输入表的windowends相等。将来,我们还可以简化joinon子句,如果窗口TVF是TUMBLE或HOP,则只包含windowstart相等。
目前(截至FLink版本1.17),窗口TVF必须与左右输入相同。这可以在未来扩展,例如,滚动窗口加入具有相同窗口大小的滑动窗口。
目前(截至FLink版本1.17),如果窗口加入在窗口化TVF之后,则窗口化TVF必须与TumbleWindows,HopWindowsorCumulateWindows一起使用,而不是Sessionwindows。
以上内容由58汽车提供。如有任何买车、用车、养车、玩车相关问题,欢迎在下方表单填写您的信息,我们将第一时间与您联系,为您提供快捷、实用、全面的解决方案。
原创文章,作者:58汽车,如若转载,请注明出处:https://car.58.com/7242505/
