Flink入门(五)--Flink算子

news2025/9/9 5:55:38

Map 


DataStream → DataStream 

一个接受一个元素并产生一个元素的函数。

示例

dataStream.map { x => x * 2 }


FlatMap 


DataStream → DataStream 

一个接受一个元素并产生零个、一个或多个元素的函数。

例如

dataStream.flatMap { str => str.split(" ") }


Filter 


DataStream → DataStream 

对于每个元素,设定一个布尔函数,并保留那些使函数返回true的元素。

例如 

dataStream.filter { _ != 0 }


KeyBy 


DataStream → KeyedStream 

逻辑上将流划分为不相交的分区。所有具有相同键的记录都被分配到同一个分区中。在内部,keyBy() 是通过哈希分区来实现的。指定键的方式有多种。

注意:没有实现hashcode()方法的POJO类和任何类型的数组都无法作为Key!!!

Reduce

KeyedStream → WindowedStream

该操作会连续地将当前元素与上一个reduce操作的结果(即最后一个reduced值)进行合并,并发出新的合并后的值。这种操作通常用于计算流数据的累积或滚动汇总。

 例如

keyedStream.reduce { _ + _ }

Window 


KeyedStream → WindowedStream 

在已经分区的KeyedStreams上可以定义窗口。窗口根据某些特性(例如,在过去5秒内到达的数据)将每个键中的数据分组。

例如 

dataStream
  .keyBy(_._1)
  .window(TumblingEventTimeWindows.of(Time.seconds(5)))

对于窗口有关的知识点可以参考我的另一篇博文

Flink入门(四) -- Flink中的窗口_flink 窗口概念 使用场景-CSDN博客

WindowAll 


DataStream → AllWindowedStream

窗口可以在常规数据流(DataStream)上定义。窗口会根据某些特性(例如,在过去5秒内到达的数据)将所有流事件进行分组。

 例如

dataStream
  .windowAll(TumblingEventTimeWindows.of(Time.seconds(5)))

Tips:在许多情况下,这是一个非并行转换。对于windowAll操作符,所有记录都将被收集到一个任务中。

Window和WindowAll的异同

特性WindowwindowAll
应用场景适用于已经分区的KeyedStream,对分区内的数据进行窗口化处理适用于未分区的DataStream,将所有流事件作为一个整体进行窗口化处理
并行度并行度是任意的,取决于后续算子的配置和KeyedStream的分区数量并行度固定为1,所有数据都被聚合到一个任务上进行处理
性能影响由于可以并行处理多个分区的数据,通常具有较好的性能由于所有数据都被聚合到一个任务上,当数据量较大时可能导致性能瓶颈
使用场景举例统计每个用户的最近5分钟内的活跃次数等需要按key分别处理的场景统计整个系统的总活跃用户数等需要对全局数据进行统计的场景,但需注意性能问题
窗口分配器与函数需要结合窗口分配器(WindowAssigner)和窗口函数(WindowFunction)来定义具体的窗口操作同样需要结合窗口分配器和窗口函数来定义窗口操作
灵活性灵活性较高,可以根据不同的key进行分区和窗口化处理灵活性较低,因为所有数据都被视为一个整体进行处理

 Window Apply

WindowedStream → DataStream

Window Apply 是一个操作,它允许你应用一个函数到整个窗口上。这意味着你可以定义一个自定义函数来处理窗口内的所有元素,而不是仅仅对每个元素进行独立的操作。这个操作的结果是产生一个新的 DataStream,其中包含了函数处理每个窗口后的结果。

 例如

windowedStream.apply { WindowFunction }

// applying an AllWindowFunction on non-keyed window stream
allWindowedStream.apply { AllWindowFunction }

 Union

与sql中union类似

DataStream* → DataStream

两个或多个数据流的联合操作会创建一个新的数据流,该数据流包含所有原始数据流中的所有元素。需要注意的是,如果你将一个数据流与自身进行联合,那么在结果数据流中,每个元素将会出现两次。(不去重不排序)

Join 

Join two data streams on a given key and a common window.

dataStream.join(otherStream)
    .where(<key selector>).equalTo(<key selector>)
    .window(TumblingEventTimeWindows.of(Time.seconds(3)))
    .apply { ... }

Interval Join

KeyedStream,KeyedStream → DataStream

例如

假设你有两个数据流:

订单流(Order Stream):包含订单的详细信息,每个订单都有一个唯一的订单ID、用户ID、订单时间戳(下单时间)和订单金额等。
支付流(Payment Stream):包含支付的详细信息,每个支付都有一个唯一的支付ID、对应的订单ID、支付时间戳和支付金额等。

你的任务是分析订单的支付情况,包括支付是否及时(例如,是否在订单下单后的几分钟内完成支付)。这里,intervalJoin 就可以派上用场了。

// this will join the two streams so that
// key1 == key2 && leftTs - 2 < rightTs < leftTs + 2
keyedStream.intervalJoin(otherKeyedStream)
    .between(Time.milliseconds(0), Time.milliseconds(20000)) 
    // lower and upper bound
    .upperBoundExclusive(true) // optional
    .lowerBoundExclusive(true) // optional
    .process(new IntervalJoinFunction() {...})

partition

  • 自定义分区

  • DataStream→DataStream 使用用户定义的分区程序为每个数据元选择目标任务。

dataStream.partitionCustom(partitioner, "someKey");
dataStream.partitionCustom(partitioner, 0);
  • 随机分区

  • DataStream→DataStream 根据均匀分布随机分配数据元。

dataStream.shuffle();
  • Rebalance (循环分区)

  • DataStream→DataStream 分区数据元循环,每个分区创建相等的负载。在存在数据倾斜时用于性能优化。

dataStream.rebalance();

   · rescaling

元素以轮询方式分区到下游操作的一个子集。这在您希望拥有这样的管道时非常有用,例如,从源的每个并行实例分发到几个映射器的子集以分散负载,但又不想触发rebalance()方法所带来的全面重新平衡。这取决于其他配置值(如TaskManager的插槽数),可能只需要本地数据传输,而不需要通过网络传输数据。

上游操作发送元素的下游操作子集取决于上游和下游操作的并行度。例如,如果上游操作有2个并行度,而下游操作有6个并行度,那么一个上游操作会将元素分发到三个下游操作,而另一个上游操作会将元素分发到另外三个下游操作。另一方面,如果下游操作有2个并行度,而上游操作有6个并行度,那么三个上游操作会将元素分发到一个下游操作,而另外三个上游操作会将元素分发到另一个下游操作。

在不同并行度不是彼此倍数的情况下,一个或多个下游操作将从上游操作接收到不同数量的输入。

dataStream.rescale()

本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若转载,请注明出处:http://www.coloradmin.cn/o/2073272.html

如若内容造成侵权/违法违规/事实不符,请联系多彩编程网进行投诉反馈,一经查实,立即删除!

相关文章

把直播间搬到工厂,淘宝直播打造卖爆新路径

又一年中秋将至&#xff0c;电商平台们再度开启了月饼生意。 8月21日&#xff0c;杭州&#xff0c;淘宝直播的主播们组成“白月光”队和“黑月牙”队&#xff0c;下工厂&#xff0c;探访体验馆&#xff0c;开始了一场“寻月之旅”。“我们米月饼的饼皮是根据南宋糕点改良而来”…

C语言小项目源码大全(60套)

C语言小项目源码大全60套 目录源码文件 目录 纯c语言迷宫源码.exe . c语言五子棋源码.exe c语言24点游戏源码.exe c语言万年历源码.exe c语言别踩白块儿(双人版)源码.exe c语言奔跑的火柴人游戏源码.exe c语言吃逗游戏源码.exe C语言超市管理系统.exe c语言对对碰游戏…

【CSP:202212-2】训练计划(Java)

题目链接 202212-2 训练计划 题目描述 求解思路 模拟&#xff1a; over表示能否按时完成所有训练项目rely[i]表示第i个项目的依赖项目编号&#xff08;每个项目最多有一个依赖项目&#xff09;days[i]用来记录第i个项目完成需要的天数allDays[i]表示加上该项目的所有前置依赖…

面向对象09:instanceof和类型转换

‌ 本节内容视频链接&#xff1a;https://www.bilibili.com/video/BV12J41137hu?p72&vd_sourceb5775c3a4ea16a5306db9c7c1c1486b5https://www.bilibili.com/video/BV12J41137hu?p72&vd_sourceb5775c3a4ea16a5306db9c7c1c1486b5 instanceof是Java中的一个二元运算符&…

浅谈【数据结构】栈和队列之队列

目录 1、队列 1.1思想 2、队列的两类 2.1顺序队列 2.2链式队列 谢谢帅气美丽且优秀的你看完我的文章还要点赞、收藏加关注 没错&#xff0c;说的就是你&#xff0c;不用再怀疑&#xff01;&#xff01;&#xff01; 希望我的文章内容能对你有帮助&#xff0c;一起努力吧&a…

MATLAB 沿任意方向分层点云(82)

MATLAB 沿任意方向分层点云(82) 一、算法介绍二、算法实现1.代码2.效果更多内容参考: MATLAB点云处理学习 一、算法介绍 沿着某个方向,将点云分割为多层,每层点云使用不同颜色进行可视化显示,具体代码和不同方向的分层效果如下: 二、算法实现 1.代码 % Load point c…

学生信息管理系统的设计与实现(包含文档、源码、sql脚本、导入视频教程)

&#x1f449;文末查看项目功能视频演示获取源码sql脚本视频导入教程视频 1 、功能描述 学生信息管理系统拥有三种角色&#xff0c;分别为学生、教师和管理员&#xff0c;功能更加完善&#xff0c;可以作为初学者参照学习课程设计。 学生&#xff1a;班级通讯录查询、个人信息…

一键生成PPT只需这一步!AI先行者全面指南

在当今快节奏的工作生活中&#xff0c;我们需要不断地准备各种报告和演示文稿。传统的PPT制作方式往往耗费大量时间和精力&#xff0c;而AI先行者的出现改变了这一切。这款强大的智能工具能够帮助您快速生成高质量的PPT&#xff0c;提高工作效率。今天&#xff0c;我们将为您带…

CLASS1:文献管理软件使用

1 文献查阅 引新(3年内)不引旧引用经典2 文献检索网站汇总 Web of Science(论文中了之后下载证明) Author Search - Web of Science Core Collection (clarivate.cn) X-MOL(查阅文献) X-MOL学术平台 计算机, 热门类期刊, - X-MOL Scidown(下载原文) Sci论文期刊检索|

zabbix监控进程、日志、主从(状态、延迟)

环境&#xff1a;rocky Linux9虚拟机四台&#xff0c;zabbix端为服务端&#xff0c;node6为客户端&#xff0c;node4为mariadb主&#xff0c;node7为mariadb从 一、zabbix监控进程 以httpd服务为例 1、客户端安装httpd [rootnode6 ~]# yum -y install httpd [rootnode6 ~]#…

微服务Gateway服务⽹关

一、Gateway服务⽹关 1.1为什么需要⽹关 Gateway⽹关是我们服务的守⻔神&#xff0c;所有微服务的统⼀⼊⼝。 ⽹关的核⼼功能特性&#xff1a; 请求路由和负载均衡&#xff1a;⼀切请求都必须先经过gateway&#xff0c;但⽹关不处理业务&#xff0c;⽽是根据某种规则&…

专利写作笔记

最近又要写专利&#xff0c;每次写专利的时候都找不到之前的专利笔记&#xff0c;这次发到网站上记录一下。 专利文件&#xff1a;1.权利要求书、2.说明书、3.说明书附图、4.说明书摘要、5.摘要附图 明确三点&#xff1a;①和现有方案的区别点&#xff08;哪个步骤不同&#x…

【02】ctf工具ECCTOOL工具的安装和使用

2.ECCTOOL工具的安装和使用 工具的介绍&#xff1a; 一款非常好用的计算ECC的工具&#xff0c;可以处理一些小数值的计算&#xff0c;点击就可以使用&#xff0c;非常方便实用&#xff0c;具体的使用方法可以参考下面图中的介绍&#xff0c;解决一定的ECC椭圆曲线的问题&…

4款在线视频压缩工具,帮你的视频文件 轻松“瘦身” 。

设备里面视频太多&#xff0c;内存不够怎么办&#xff1f;视频文件太大不好传输怎么办&#xff1f;视频文件大小受规则限制怎么办&#xff1f; 别担心&#xff01;有了这4款视频压缩软件&#xff0c;轻轻松松帮你搞定这些问题。 1、福昕视频高效压缩 直通车&#xff1a;www.f…

进制转换计算幸运数出现次数(华为od机考题)

一、题目 1.原题 有位客人来自异国&#xff0c;在该国使用m进制计数。 该客人有个幸运数字n(n<m)&#xff0c;每次购物时&#xff0c; 其总是喜欢计算本次支付的花费(折算为异国的价格后)中存在多少幸运数字。 问&#xff1a;当其购买一个在我国价值k的产品时&#xff0c;…

AI在医学领域:GluFormer一种可泛化的连续血糖监测数据分析基础模型

糖尿病是一种全球性的健康挑战&#xff0c;影响着各个年龄段和不同地理区域的人群。根据最新数据&#xff0c;全球糖尿病患者人数已超过5亿&#xff0c;且每年以惊人的速度增长&#xff0c;相关的医疗费用也居高不下。2型糖尿病&#xff08;T2DM&#xff09;作为最主要的糖尿病…

lit-llama代码笔记--LLaMA Model

代码来自&#xff1a;lit-llama modelscope模型下载 &#xff1a;llama-7b 下载后的模型需要转换为lit-llama使用的格式&#xff0c;详见 howto 文件夹下的 download_weights.md 文中代码为了方便说明&#xff0c;删减了一些内容&#xff0c;详细代码请查看源码。 generate …

u盘突然说要格式化才能访问?如何跳过格式化打开U盘

在日常使用U盘的过程中&#xff0c;有时我们会突然遇到U盘无法直接访问&#xff0c;系统提示需要格式化才能继续使用的情况。这往往让人措手不及&#xff0c;尤其是当U盘中存储着重要数据时。面对这样的困境&#xff0c;许多用户可能会感到焦虑和无助。然而&#xff0c;不必过于…

SQLserver中的触发器和存储过程

在 SQL Server 中&#xff0c;触发器是一种特殊的存储过程&#xff0c;它在指定的数据库表上发生特定的数据修改事件时自动执行。触发器可以用于执行各种任务&#xff0c;如数据验证、数据审计、自动更新相关表等。 触发器的类型 SQL Server 支持以下几种类型的触发器&#x…

如何构建基于Java SpringBoot的保险业务管理与数据分析系统

✍✍计算机编程指导师 ⭐⭐个人介绍&#xff1a;自己非常喜欢研究技术问题&#xff01;专业做Java、Python、微信小程序、安卓、大数据、爬虫、Golang、大屏等实战项目。 ⛽⛽实战项目&#xff1a;有源码或者技术上的问题欢迎在评论区一起讨论交流&#xff01; ⚡⚡ Java实战 |…