聊聊flink Table的Group Windows

本文旨在探讨flink table的group windows

聊聊flink Table的Group Windows

Table table = input    .window([Window w].as("w"))  // 定义窗口并为其赋予别名 w    .groupBy("w")  // 按窗口 w 分组表    .select("b.sum");  // 聚合Table table = input    .window([Window w].as("w"))  // 定义窗口并为其赋予别名 w    .groupBy("w, a")  // 按属性 a 和窗口 w 分组表    .select("a, b.sum");  // 聚合Table table = input    .window([Window w].as("w"))  // 定义窗口并为其赋予别名 w    .groupBy("w, a")  // 按属性 a 和窗口 w 分组表    .select("a, w.start, w.end, w.rowtime, b.count"); // 聚合并添加窗口的开始、结束和行时间戳

窗口操作可以为Window设置别名,并在groupBy及select中引用该别名。窗口具有start、end和rowtime属性,其中start和rowtime是包含的,而end是排外的。

Tumbling Windows:

// 事件时间的Tumbling窗口.window(Tumble.over("10.minutes").on("rowtime").as("w"));// 处理时间的Tumbling窗口(假设有一个处理时间属性 "proctime").window(Tumble.over("10.minutes").on("proctime").as("w"));// 基于行数的Tumbling窗口(假设有一个处理时间属性 "proctime").window(Tumble.over("10.rows").on("proctime").as("w"));

Tumbling Windows按照固定窗口大小移动,因此窗口之间不重叠;over方法用于指定窗口大小;窗口大小可以基于事件时间、处理时间或行数来定义。

Sliding Windows:

// 事件时间的Sliding窗口.window(Slide.over("10.minutes").every("5.minutes").on("rowtime").as("w"));// 处理时间的Sliding窗口(假设有一个处理时间属性 "proctime").window(Slide.over("10.minutes").every("5.minutes").on("proctime").as("w"));// 基于行数的Sliding窗口(假设有一个处理时间属性 "proctime").window(Slide.over("10.rows").every("5.rows").on("proctime").as("w"));

当滑动间隔小于窗口大小时,Sliding Windows会导致窗口重叠,因此行可能属于多个窗口;over方法用于指定窗口大小,窗口大小可以基于事件时间、处理时间或行数来定义;every方法用于指定滑动间隔。

Session Windows:

// 事件时间的Session窗口.window(Session.withGap("10.minutes").on("rowtime").as("w"));// 处理时间的Session窗口(假设有一个处理时间属性 "proctime").window(Session.withGap("10.minutes").on("proctime").as("w"));

Session Windows没有固定的窗口大小,它基于非活动时间的长度来关闭窗口,withGap方法用于指定两个窗口之间的间隔,作为时间间隔;Session Windows只能使用事件时间或处理时间。

Table类提供了window操作,接收Window参数,并创建WindowedTable对象。

class Table(    private[flink] val tableEnv: TableEnvironment,    private[flink] val logicalPlan: LogicalNode) {  //......  def window(window: Window): WindowedTable = {    new WindowedTable(this, window)  }  //......}

WindowedTable类仅提供groupBy操作,groupBy可以接收String类型的参数,也可以接收Expression类型的参数;String类型的参数会被转换为Expression类型,最终调用的是Expression类型的groupBy方法;如果groupBy操作除了窗口之外没有其他属性,则其并行度为1,只会在单个任务上执行;groupBy方法创建WindowGroupedTable对象。

class WindowedTable(    private[flink] val table: Table,    private[flink] val window: Window) {  def groupBy(fields: Expression*): WindowGroupedTable = {    val fieldsWithoutWindow = fields.filterNot(window.alias.equals(_))    if (fields.size != fieldsWithoutWindow.size + 1) {      throw new ValidationException("GroupBy must contain exactly one window alias.")    }    new WindowGroupedTable(table, fieldsWithoutWindow, window)  }  def groupBy(fields: String): WindowGroupedTable = {    val fieldsExpr = ExpressionParser.parseExpressionList(fields)    groupBy(fieldsExpr: _*)  }}

WindowGroupedTable类仅提供select操作,select可以接收String类型的参数,也可以接收Expression类型的参数;String类型的参数会被转换为Expression类型,最终调用的是Expression类型的select方法;select方法创建新的Table对象,其Project操作的子节点为WindowAggregate

class WindowGroupedTable(    private[flink] val table: Table,    private[flink] val groupKeys: Seq[Expression],    private[flink] val window: Window) {  def select(fields: Expression*): Table = {    val expandedFields = expandProjectList(fields, table.logicalPlan, table.tableEnv)    val (aggNames, propNames) = extractAggregationsAndProperties(expandedFields, table.tableEnv)    val projectsOnAgg = replaceAggregationsAndProperties(      expandedFields, table.tableEnv, aggNames, propNames)    val projectFields = extractFieldReferences(expandedFields ++ groupKeys :+ window.timeField)    new Table(table.tableEnv,      Project(        projectsOnAgg,        WindowAggregate(          groupKeys,          window.toLogicalWindow,          propNames.map(a => Alias(a._1, a._2)).toSeq,          aggNames.map(a => Alias(a._1, a._2)).toSeq,          Project(projectFields, table.logicalPlan).validate(table.tableEnv)        ).validate(table.tableEnv),        // required for proper resolution of the time attribute in multi-windows        explicitAlias = true      ).validate(table.tableEnv))  }  def select(fields: String): Table = {    val fieldExprs = ExpressionParser.parseExpressionList(fields)    //get the correct expression for AggFunctionCall    val withResolvedAggFunctionCall = fieldExprs.map(replaceAggFunctionCall(_, table.tableEnv))    select(withResolvedAggFunctionCall: _*)  }}

总结:窗口操作可以为Window设置别名,并在groupBy及select中引用该别名。窗口具有start、end和rowtime属性,其中start和rowtime是包含的,而end是排外的。Tumbling Windows按固定窗口大小移动,不重叠;Sliding Windows在滑动间隔小于窗口大小的情况下会重叠;Session Windows基于非活动时间关闭窗口。Table类提供window操作,创建WindowedTable;WindowedTable提供groupBy操作,创建WindowGroupedTable;WindowGroupedTable提供select操作,创建新的Table,其Project操作的子节点为WindowAggregate。

以上就是聊聊flink Table的Group Windows的详细内容,更多请关注创想鸟其它相关文章!

版权声明:本文内容由互联网用户自发贡献,该文观点仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。
如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 chuangxiangniao@163.com 举报,一经查实,本站将立刻删除。
发布者:程序猿,转转请注明出处:https://www.chuangxiangniao.com/p/174829.html

(0)
打赏 微信扫一扫 微信扫一扫 支付宝扫一扫 支付宝扫一扫
品牌手机壳排行榜前十名推荐(最实用的手机壳品牌)
上一篇 2026年8月31日 03:40:17
Tailwind CSS hocus 变体失效:为什么按钮焦点状态下样式不生效?
下一篇 2026年8月31日 03:44:34

相关推荐

  • win10如何解决“Geolocation Service”无法定位或持续运行的问题_修复定位服务异常的方法

    首先通过服务管理器启动Geolocation Service,若无效则修复或删除注册表中lfsvc下的TriggerInfo项以重置服务状态,最后检查组策略是否禁用定位服务并进行调整。 如果您尝试在Windows 10中启用定位功能,但发现“Geolocation Service”服务无法启动或持续…

    2026年9月10日
    000
  • edge浏览器如何设置在新标签页中打开链接_Edge浏览器设置新标签页打开链接方法

    可通过鼠标中键点击、按住Ctrl键左击、右键选择“在新标签页中打开”、安装扩展程序或使用组策略五种方法实现链接在新标签页打开。 如果您在浏览网页时希望链接能够在新标签页中打开,从而避免覆盖当前页面内容,可以通过调整 Microsoft Edge 浏览器的设置或使用特定操作方式实现。以下是几种可行的方…

    2026年9月10日
    200
  • Laravel如何调度定时任务_自动化任务调度配置

    Laravel的定时任务调度通过将Cron配置集中到代码中,解决了传统方式的分散、难维护问题。核心在于创建Artisan命令并在app/Console/Kernel.php的schedule方法中定义调度逻辑,如使用dailyAt()设置执行时间,withoutOverlapping()防止重复执行…

    2026年9月10日
    100
  • Linux系统如何防止提权攻击_Linux防止提权攻击的防护措施

    防范提权攻击需坚持最小权限原则,合理配置用户权限与sudo策略,加固文件目录权限,定期检查SUID/SGID滥用,及时更新系统补丁,启用SELinux/AppArmor等安全模块,并通过auditd和日志监控实现异常行为检测,结合权限管理、系统加固与实时监控形成综合防护体系。 防止提权攻击是Linu…

    2026年9月10日
    100
  • thinkphp如何实现文件上传功能

    ThinkPHP 6 实现文件上传需创建上传目录并设置可写权限,前端表单使用 multipart/form-data 编码,控制器通过 Request::file() 获取文件,利用 Filesystem 组件的 putFile() 方法自动重命名并保存至 public/storage 目录,支持 …

    2026年9月10日
    100
  • windows怎么禁用搜索索引_Windows搜索索引功能关闭方法

    关闭Windows搜索索引可降低系统资源占用,通过服务管理器禁用Windows Search服务、使用组策略编辑器或注册表修改,并删除索引数据库文件以释放空间。 如果您发现Windows系统的搜索功能响应缓慢或占用过多系统资源,可能是由于后台的搜索索引服务持续运行所致。禁用搜索索引可以有效降低磁盘和…

    2026年9月10日
    100
  • OPPO Find X8 Pro拍照慢怎么办 OPPO Find X8 Pro相机优化

    拍照慢多因设置不当或系统卡顿,先检查相机设置并优化使用习惯。1. 避免常开大师模式,日常用自动模式提升响应速度;2. 关闭AI超清像素、AI千里长焦等耗时功能以减少计算延迟;3. 利用侧边快捷键双击秒开相机、单击拍照、长按连拍,提升抓拍效率;4. 精简取景界面,关闭水印和网格线使操作更直观;5. 重…

    2026年9月10日
    100
  • 如何在mysql中分析复制日志错误

    首先检查复制状态,使用SHOW SLAVE STATUSG查看Slave_IO_Running和Slave_SQL_Running状态及Last_Error信息;再分析错误日志文件hostname.err中与“[ERROR]”或“Replication”相关的记录;最后根据主键冲突、GTID不一致、…

    2026年9月10日
    200
  • 如何为VSCode配置Rust开发环境?

    首先安装Rust工具链并验证rustc与cargo版本,接着在VSCode中安装Rust Analyzer和CodeLLDB插件,配置settings.json实现保存时自动格式化与clippy检查,最后通过launch.json设置调试环境,确保项目可运行调试。 为 VSCode 配置 Rust …

    2026年9月10日
    1000
  • 腾讯ima公布2.0版本,开启任务模式内测,可通过Agent能力生成报告和播客

    借助智能体技术实现自主任务规划,新版ima推出“任务模式”,支持自动生成报告与播客内容。用户可在首页或知识库中以自然语言形式发起任务请求,并附加文档、图片、音频、网页链接、笔记及知识库资料等多种素材,为大模型执行任务提供详实的“参考资料”。在生成播客时,系统还允许用户自定义对谈人数、选择音色等参数,…

    2026年9月10日
    200
  • iPhone 17 能用挂绳了,梦回 20 年前?

    iPhone 17 能用挂绳了,梦回 20 年前?iPhone 17 能用挂绳了,梦回 20 年前?iPhone 17 能用挂绳了,梦回 20 年前?iPhone 17 能用挂绳了,梦回 20 年前?

    还记得 2006 年左右,手机挂绳风靡一时,像诺基亚、索尼爱立信、摩托罗拉等品牌的机型几乎都配备了挂绳孔,方便用户挂上装饰物,既防丢又能展现个性。 尤其是在非主流文化盛行的年代,搭配跑马灯手机、杀马特发型和服饰,一条挂绳瞬间让整体造型更具“潮流感”。 年代久远,没找到更合适的图 不过随着技术进步,厂…

    2026年9月10日 用户投稿
    200
  • Java 函数中的布尔类型返回值

    本文详细讲解了如何在 java 中创建和使用返回布尔类型的函数,以判断一个数是否为质数为例,展示了如何避免变量初始化问题,并提供了优化后的代码示例,帮助开发者编写更简洁高效的 java 代码。 在 Java 编程中,经常需要编写函数来执行特定的任务并返回一个结果。布尔类型(boolean)是一种常用…

    2026年9月10日
    100
  • Java双向路径搜索实现详解与路径构建指南

    本文旨在帮助开发者理解和实现Java中的双向路径搜索算法。我们将深入探讨算法的核心思想,并针对常见的实现错误进行分析。通过改进代码逻辑,我们将展示如何构建完整的从起始点到终点的路径,确保算法的正确性和效率。 双向路径搜索的核心思想 双向路径搜索是一种优化搜索算法,它同时从起始节点和目标节点开始搜索,…

    2026年9月10日
    100
  • Windows10自动更新怎么彻底关闭_Windows10自动更新关闭方法

    可通过组策略、注册表、服务禁用等方法彻底关闭Windows 10自动更新。一、专业版用户使用gpedit.msc进入“配置自动更新”设为禁用,并启用“删除所有更新功能访问权限”。二、家庭版用户通过regedit在HKEY_LOCAL_MACHINESOFTWAREPoliciesMicrosoftW…

    2026年9月10日
    100
  • thinkphp多应用模式下公共模块如何创建

    创建公共模块需在根目录下建立common目录并配置PSR-4自动加载,通过命名空间在多应用间共享模型、服务与中间件,实现代码复用。 在 ThinkPHP 多应用模式下,公共模块的创建主要是为了解决多个应用之间共享模型、服务、工具类或配置的问题。通过合理组织目录结构和自动加载机制,可以实现代码复用,避…

    2026年9月10日
    100
  • 王腾:REDMI Note 15 Pro 新机本月发布 顶级标准打造

    8 月 12 日,redmi 品牌总经理王腾在社交媒体上分享了两大重磅消息。其一,redmi note 系列已成功进入全球超过一百个国家,成为 2025 年上半年 175-499 美元价格区间内全球最畅销的国产手机。 REDMI Note 14 系列 REDMI Note 系列素来被称为 &#822…

    2026年9月10日
    000
  • REDMI最强旗舰来了!博主称REDMI K90全系应该都会涨价

    10月23日消息,博主数码闲聊站表示,redmi k90系列应该是全系涨价,具体涨多少不知道。 此前REDMI产品经理笋寸称,今年内存涨价很多,米粉对K90系列定价得有一定的心理准备。 据悉,上一代旗舰K80系列标准版起售价是2499元,K80 Pro起售价是3699元,涨价意味着K90系列起售价至…

    2026年9月10日
    000
  • AIGC查重官网入口 知网免费检测链接

    知网官方不提供免费AIGC查重服务,个人用户需通过https://cx.cnki.net官网付费检测,或经由学校等机构系统获取免费权限,谨防非官方“免费入口”风险。 ☞☞☞AI 智能聊天, 问答助手, AI 智能搜索, 免费无限量使用 DeepSeek R1 模型☜☜☜ 目前知网官方不提供免费的AI…

    2026年9月10日
    100
  • VSCode注释技巧:生成文档与注解

    掌握VSCode%ignore_a_1%技巧可提升代码可读性与开发效率。1. 使用JSDoc添加函数说明,支持智能提示;2. 快捷键Ctrl/Cmd+/快速切换行注释,输入/**自动生成块注释;3. 配合”Document This”插件一键生成JSDoc模板;4. 利用js…

    2026年9月10日
    200
  • 如何在安装mysql时设置默认存储引擎

    在MySQL配置文件的[mysqld]段落中添加default-storage-engine=InnoDB,2. 初始化时可通过命令指定默认引擎,3. 启动后执行SHOW VARIABLES验证设置,创建表并用SHOW CREATE TABLE确认引擎类型是否生效。 在安装 MySQL 时设置默认存…

    2026年9月10日
    100

发表回复

登录后才能评论
关注微信