【大数据技术攻关专题】「Apache-Flink零基础入门」手把手+零基础带你玩转大数据流式处理引擎Flink(基础加强+运行原理)

手把手+零基础带你玩转大数据流式处理引擎Flink(运行机制+原理加深)

  • 前提介绍
  • 运行Flink应用
    • 运行机制
      • Flink的两大核心组件
        • JobManager
        • TaskManager
          • TaskSlot
      • Flink分层架构
        • Stateful Stream Processing
        • DataStream和DataSet
          • DataStream(数据流)
            • 特点
          • DataSet(数据集):
            • 特点
        • Table SQL
      • DataStream API 编程模型
        • 批处理
        • 流处理
        • 流式处理系统
          • 流处理系统的特点
          • DAG图
          • 物理模型并行计算

前提介绍

关于Flink服务的搭建与部署,由于其涉及诸多实战操作而理论部分相对较少,小编打算采用一个独立的版本和环境来进行详尽的实战讲解。考虑到文字描述可能无法充分展现操作的细节和流程,我们决定以视频的形式进行分析和介绍。因此,在本文中,我们将暂时不涉及具体的搭建和部署步骤。

为确保大家能够更直观地掌握Flink服务的搭建与部署技巧,我们将专注于制作高质量的教学视频。后续,我们还会编写一篇与视频内容相辅相成的辅助教材,以帮助大家更好地理解和巩固所学知识。目前,我们的首要任务是录制部署视频,敬请期待!

运行Flink应用

在运行Flink应用之前,深入了解其运行时组件是必不可少的环节,因为这些组件的配置直接关系到应用的性能和稳定性。

运行机制

在Flink的分布式计算框架中,Task是其资源调度的最小单位。正确理解和配置这些组件,对于确保Flink应用的稳定运行和高效性能至关重要。如下图所示:
在这里插入图片描述
上图所展示的是一个使用DataStream API编写的数据处理程序。在图中,我们可以清晰地看到,那些无法被串联在一起的Operator被分隔到了不同的Task中。

Flink的两大核心组件

Flink作为流处理领域的佼佼者,其高效稳定的运行离不开两大核心组件的密切协作:JobManager和TaskManager。它们各司其职,共同支撑着整个Flink运作体系的顺畅运行,如下图所示:
在这里插入图片描述

JobManager

JobManager(也被称作JobMaster)是Flink作业执行的“大脑”,负责协调Task的分布式执行。具体来说,它会负责调度Task,确保它们在集群中的各个节点上按计划执行。

负责协调创建Checkpoint,这是Flink的容错机制之一,用于在作业发生故障时能够恢复到之前的状态。当Job发生failover时,JobManager会协调各个Task从最近的Checkpoint恢复,确保作业的持续执行。

TaskManager

TaskManager(也被称作Worker)则是Flink作业执行的“手脚”,负责具体执行Dataflow中的Tasks。它会分配内存Buffer,确保数据在各个Task之间高效传递。

执行Data Stream的处理逻辑,包括数据的接收、处理和发送等。通过多个TaskManager的并行执行,Flink能够实现大规模数据的实时处理和分析。

下面时JobManager和TaskManager连个核心组件的整体合作处理的架构图:
在这里插入图片描述

TaskSlot

从下面的图中可以看出来Task Slot是TaskManager中的最小资源分配单位,它决定了TaskManager能够支持的并发Task处理数量。
在这里插入图片描述
一个TaskManager中的Task Slot数量直接影响到其并发处理能力和资源利用率。通过合理配置Task Slot的数量,可以根据实际需求调整TaskManager的工作负载,从而实现更高效的任务处理和资源利用。

Flink分层架构

在这里插入图片描述

Stateful Stream Processing

Flink的分层架构分析,位于架构的最底层核心部分,ProcessFunction扮演着实现Flink Core API基础逻辑的关键角色。它提供了直接操作和处理流数据流的底层接口,使得开发者能够基于此构建出高度定制化的组件或功能模块,例如通过巧妙利用其内置的定时机制进行特定条件下的数据匹配与缓存。

尽管ProcessFunction带来了无可比拟的灵活性,允许对数据流处理过程进行细粒度控制,但同时也意味着开发工作相对复杂,需要对Flink的工作原理和并行计算有深入理解才能更好地驾驭这一强大工具。

DataStream和DataSet

DataStream 和 DataSet 是两个核心概念。它们是 Flink 中用于处理数据的两种不同的抽象。
在这里插入图片描述

  • DataStream:适用于处理连续的实时数据流,提供了丰富的流处理操作符和函数,可以实现实时流处理的需求;
  • DataSet:适用于处理有限的离线数据集,提供了丰富的批处理操作符和函数,可以实现离线数据处理的需求。
DataStream(数据流)

DataStream 是 Flink 中处理连续流数据的抽象。它表示无限的数据流,可以是来自消息队列、日志文件、传感器等源的实时数据。

特点
  • 有序的、可变长度的数据记录序列,每个记录都包含一个或多个字段。每个记录都有一个时间戳,用于指示记录的时间顺序。
  • 丰富的操作符和函数,可以对数据流进行转换、过滤、聚合等操作。可以通过窗口操作来处理有限大小的数据窗口,也可以进行流处理的时间语义控制。
  • 基于事件时间(Event Time)或处理时间(Processing Time)进行处理的,可以实现事件驱动的流处理。
DataSet(数据集):

DataSet 是 Flink 中处理有限数据集的抽象。它表示有限的、静态的数据集合,可以是来自文件、数据库、批处理作业等离线数据。

特点
  • 不可变的、有限长度的数据集合,每个数据集合由一组记录组成,每个记录都包含一个或多个字段。
  • 丰富的操作符和函数,可以对数据集进行转换、过滤、聚合等操作。可以通过分组、排序、连接等操作来处理数据集。
  • 基于批处理模式进行处理的,适用于离线数据处理和批处理作业。
Table SQL
  • SQL 是基于 Table 的,因此在使用 SQL 之前需要创建一个 Table 环境。
  • 不同类型的 Table 需要使用相应的 Table 环境进行构建。
  • Table 可以与 DataStream 或 DataSet 相互转换,这使得我们可以在流处理和批处理之间无缝切换。
  • Streaming SQL 与存储的 SQL 有所不同,它会被转化为流式执行计划,以实现实时流处理的需求。

后面的章节会针对性详细介绍,此处大概了解就可以了。

DataStream API 编程模型

流处理和批处理是大数据处理中的两个核心概念,它们从不同的角度对数据进行处理。它们的关系可以类比于 Java 中的 ArrayList 中的元素,可以通过下标直接访问,也可以通过迭代器进行访问。

批处理

批处理是对有限的静态数据集进行处理的方式。它以批量的方式处理数据,数据是一次性加载并进行处理。批处理适用于离线数据处理和批量作业,如数据清洗、数据分析等。在批处理中,数据被视为有限的数据集合,可以通过分组、排序、连接等操作进行处理。

流处理

流处理是对连续的实时数据流进行处理的方式。它以事件驱动的方式处理数据,数据是逐个到达的,并且可以立即进行处理。流处理适用于实时性要求较高的场景,如实时监控、实时分析等。在流处理中,数据被视为无限的流,可以通过窗口操作来处理有限大小的数据窗口。

流式处理系统

流处理系统具有许多独特的特点。通常情况下,由于需要处理无限数据集,流处理系统采用一种数据驱动的处理方式。它会预先设置一些算子,并在数据到达时对数据进行处理。

在这里插入图片描述

流处理系统的特点
  1. 实时处理:流处理系统能够实时处理连续的数据流,数据到达后立即进行处理,实现实时性要求较高的应用场景。
  2. 无限数据集:流处理系统能够处理无限的数据流,不受数据大小的限制。它能够处理持续不断产生的数据,而不需要等待所有数据都可用。
  3. 数据驱动:流处理系统是以数据为驱动的,即数据到达时才进行处理。系统会根据数据的到达情况来触发相应的处理操作,而不是按照固定的时间间隔进行处理。
DAG图

为了表达复杂的计算逻辑,包括 Flink 在内的分布式流处理引擎一般采用DAG图来表示整个计算逻辑,其中 DAG 图中的每一个点就代表一个基本的逻辑单元,也就是算子。

由于计算逻辑被组织成有向图,数据会按照边的方向,从一些特殊的 Source 节点流入系统,然后通过网络传输、本地传输等不同的数据传输方式在算子之间进行发送和处理,最后会通过另外一些特殊的 Sink 节点将计算结果发送到某个外部系统或数据库中,下面是执行计划的DAG逻辑图
在这里插入图片描述

物理模型并行计算

在物理模型中,我们根据计算逻辑的需求,通过系统自动优化或人为指定的方式将计算工作分布到不同的实例中。只有当算子实例分布到不同的进程上时,才需要通过网络进行数据传输。而在同一进程中的多个实例之间的数据传输通常是不需要通过网络的。

通过将计算工作分布到不同的实例中,可以实现并行计算和分布式处理,以提高整体的计算性能和吞吐量。在分布式流处理引擎中,系统会根据算子的并行度和资源配置,将算子实例分布到不同的计算节点上,使得每个实例可以独立地处理数据。

对于实际的分布式流处理引擎,它们的实际运行时物理模型要更复杂一些,这是由于每个算子都可能有多个实例。如下图所示:
在这里插入图片描述
在实际的分布式流处理引擎中,物理模型比逻辑模型更加复杂。这种复杂性是由于分布式流处理引擎的并行性和分布式计算的特性所导致的。为了实现高吞吐量和低延迟的数据处理,引擎会将算子实例分布在多个计算节点上,并通过网络进行数据交换和通信。

例如,图中的算子 A 作为数据源有两个实例,而中间算子 C 也有两个实例。在逻辑模型中,A 和 B 是 C 的上游节点,但在物理模型中,C 的每个实例与 A 和 B 的每个实例之间都可能存在数据交换。

当算子实例分布到不同的进程上时,数据传输就会发生。这时,需要通过网络进行数据的传输和交换。而在同一进程中的多个实例之间,数据传输通常是通过共享内存或进程间通信的方式进行,而不需要通过网络。
在这里插入图片描述

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

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

相关文章

GPIO中断

1.EXTI简介 EXTI是External Interrupt的缩写,指外部中断。在嵌入式系统中,外部中断是一种用于处理外部事件的机制。当外部事件发生时(比如按下按钮、传感器信号变化等),外部中断可以立即打断正在执行的程序&#xff0…

大红喜庆版UI猜灯谜小程序源码/猜字谜微信小程序源码

今天给大家带来一款UI比较喜庆的猜灯谜小程序,大家看演示图的时候当然也是可以看得到那界面是多么的喜庆,而且新的一年也很快就来了,所以种种的界面可能都比较往喜庆方面去变吧。 这款小程序搭建是免服务器和域名的,只需要使用微信开发者工具…

Linux一键部署telegraf 实现Grafana Linux 图形展示

influxd2前言 influxd2 是 InfluxDB 2.x 版本的后台进程,是一个开源的时序数据库平台,用于存储、查询和可视化时间序列数据。它提供了一个强大的查询语言和 API,可以快速而轻松地处理大量的高性能时序数据。 telegraf 是一个开源的代理程序,它可以收集、处理和传输各种不…

Linux开发工具

前言:哈喽小伙伴们,经过前边的学习我们已经掌握了Linux的基本指令和权限,相信大家学完这些之后都会对Linux有一个更加深入的认识,但是Linux的学习可以说是从现在才刚刚开始。 这篇文章,我们将讲解若干个Linux的开发工…

Java基础数据结构之Map和Set

Map和Set接口 1.Set集合:独特性与无序性 Set是Java集合框架中的一种,它代表着一组无序且独特的元素。这意味着Set中的元素不会重复,且没有特定的顺序。Set接口有多个实现类,如HashSet、LinkedHashSet和TreeSet。 2.Map集合&…

Redis核心技术与实战【学习笔记】 - 19.Pika:基于SSD实现大容量“Redis”

前言 随着业务数据的增加(比如电商业务中,随着用户规模和商品数量的增加),就需要 Redis 能保存更多的数据。你可能会想到使用 Redis 切片集群,把数据分散保存到不同的实例上。但是这样做的话,如果要保存的…

利用牛顿方法求解非线性方程(MatLab)

一、算法原理 1. 牛顿方法的算法原理 牛顿方法(Newton’s Method),也称为牛顿-拉弗森方法,是一种用于数值求解非线性方程的迭代方法。其基本思想是通过不断迭代来逼近方程的根,具体原理如下: 输入&#…

PCB笔记(二十三):allegro 标注长宽(一般用于测量板宽)时如何显示双单位

步骤:首先选择标注工具,然后右键→Parameters,在弹出来的窗口中√上如下图二所示选项 最终要达到显示单位的效果的话,需要在Text项键入%v%u。 今天就记录到这里啦O

Leetcode206:反转链表

一、题目 给你单链表的头节点 head ,请你反转链表,并返回反转后的链表 示例: 输入:head [1,2,3,4,5] 输出:[5,4,3,2,1]输入:head [1,2] 输出:[2,1]输入:head [] 输出&#xff1…

[ESP32 IDF]web server

目录 通过web server控制LED 核心原理解析 分区表 web server的使用 错误Header fields are too long的解决 通过web server控制LED 通过网页控制LED灯的亮灭,一般的ESP32开发板都可以实现,下面这篇文章是国外开发者提供的一个通过web server控制…

13.2K Star,12306 抢票助手帮你回家

Hi,骚年,我是大 G,公众号「GitHub指北」会推荐 GitHub 上有趣有用的项目,一分钟 get 一个优秀的开源项目,挖掘开源的价值,欢迎关注。 马上过年了,今年你在哪里过年?回老家吗&#x…

前端入门第一天

目录 HTML超文本标记语言——Hyper Text Markup Language 标签语法: 双标签: 单标签——只有开始标签,没有结束标签 基本骨架: 标签的关系: 注释: 标题标签:(新闻标题、文章标题、网页区域…

探索设计模式的魅力:为什么你应该了解装饰器模式-代码优化与重构的秘诀

设计模式专栏:http://t.csdnimg.cn/nolNS 开篇 在一个常常需要在不破坏封装的前提下扩展对象功能的编程世界,有一个模式悄无声息地成为了高级编程技术的隐形冠军。我们日复一日地享受着它带来的便利,却往往对其背后的复杂性视而不见。它是怎样…

项目安全-----加密算法实现

目录 对称加密算法 AES (ECB模式) AES(CBC 模式)。 非对称加密 对称加密算法 对称加密算法,是使用相同的密钥进行加密和解密。使用对称加密算法来加密双方的通信的话,双方需要先约定一个密钥,加密方才能加密&#…

C++ 动态规划 线性DP 最长共同子序列

给定两个长度分别为 N 和 M 的字符串 A 和 B ,求既是 A 的子序列又是 B 的子序列的字符串长度最长是多少。 输入格式 第一行包含两个整数 N 和 M 。 第二行包含一个长度为 N 的字符串,表示字符串 A 。 第三行包含一个长度为 M 的字符串,表…

Servlet+Ajax实现对数据的列表展示(极简入门)

目录 1.准备工作 1.数据库源(这里以Mysql为例) 2.映射实体类 3.模拟三层架构(Dao、Service、Controller) Dao接口 Dao实现 Service实现(这里省略Service接口) Controller层(或叫Servlet层…

【blender插件】(1)快速开始

特性 blender的python API有如下特性: 编辑用户界面可以编辑的任意数据(场景,网格,粒子等)。修改用户首选项、键映射和主题。运行自己的配置运行工具。创建用户界面元素,如菜单、标题和面板。创建新的工具。场景交互式工具。创建与Blender集成的新渲染引擎。修改模型的数据…

LeetCode.1686. 石子游戏 VI

题目 题目链接 分析 本题采取贪心的策略 我们先假设只有两个石头a,b, 对于 Alice 价值分别为 a1,a2, 对于 Bob 价值而言价值分别是 b1,b2 第一种方案是 Alice取第一个,Bob 取第二个,Alice与Bob的价值差是 c1 a1 - b1&#xf…

【Crypto | CTF】BUUCTF rsarsa1

天命:第二题RSA解密啦,这题比较正宗 先来看看RSA算法 这道题给出了 p,q,E,就是给了两个质数和公钥 有这三个东西,那就可以得出私钥了 最后把私钥和质数放进去解密即可得到解密后的明文 from gmpy2 impor…