【大数据】流处理基础概念(一):Dataflow 编程基础、并行流处理

流处理基础概念(一):Dataflow 编程基础、并行流处理

  • 1.Dataflow 编程基础
    • 1.1 Dataflow 图
    • 1.2 数据并行和任务并行
    • 1.3 数据交换策略
  • 2.并行流处理
    • 2.1 延迟与吞吐
      • 2.1.1 延迟
      • 2.1.2 吞吐
      • 2.1.3 延迟与吞吐
    • 2.2 数据流上的操作
      • 2.2.1 数据接入和数据输出
      • 2.2.2 转换操作
      • 2.2.3 滚动聚合
      • 2.2.4 窗口操作
        • 2.2.4.1 滚动窗口
        • 2.2.4.2 滑动窗口
        • 2.2.4.3 会话窗口
        • 2.2.4.4 小结

1.Dataflow 编程基础

1.1 Dataflow 图

Dataflow 程序描述了数据如何在不同操作之间流动。Dataflow 程序通常表示为 有向图。图中 顶点 称为 算子,表示计算;而 表示 数据依赖关系算子是 Dataflow 程序的基本功能单元,它们从输入获取数据对其进行计算,然后产生数据并发往输出,以供后续处理。没有输入端的算子称为 数据源data sources),没有输出端的算子称为 数据汇data sinks)。一个 Dataflow 图至少要有一个数据源和一个数据汇。

在这里插入图片描述
上图中的 Dataflow 图被称作 逻辑图,因为它们表达了高层视角下的计算逻辑。为了执行 Dataflow 程序,需要将逻辑图转化为 物理 Dataflow 图。后者会指定程序的执行细节,例如,当我们使用分布式处理引擎时,每个算子可能会在不同物理机器上运行多个并行任务。在逻辑 Dataflow 图中,顶点代表算子。在物理 Dataflow 图中,顶点代表任务。如下图所示,“抽取主题标签” 和 “计数” 算子都包含 2 个并行算子任务,每个任务负责计算一部分输入数据。
在这里插入图片描述

1.2 数据并行和任务并行

Dataflow 图的并行性可以通过多种方式加以利用。

首先你可以将输入数据分组,让同一操作的多个任务并行执行在不同数据子集上,这种并行称为 数据并行。数据并行非常有用,因为它能够将计算负载分配到多个节点上,从而允许处理大规模的数据。

再者,你可以让不同算子的任务(基于相同或不同的数据)并行计算,这种并行称为任务并行。通过任务并行,可以更好的利用集群的计算资源。

1.3 数据交换策略

数据交换策略定义了如何将数据项分配给物理 Dataflow 图中的不同任务。这些策略可以由执行引擎根据算子的语义自动选择,也可以由 Dataflow 编程人员显示指定。常见的数据交换策略如下图所示:转发广播基于键值随机

在这里插入图片描述

2.并行流处理

在此,我们给出数据流的定义:数据流是一个可能无限的事件序列。数据流中的事件可以表示 监控数据传感器测量值信用卡交易气象站观测数据在线用户交互,以及 网络搜索 等。

2.1 延迟与吞吐

对于批处理应用而言,我们通常会关心作业的总执行时间,或者说处理引擎读取输入、执行计算、写回结果总共需要多长时间。但由于流式应用会持续执行且输入可能是无限的,所以在数据流处理中,没有总执行时间的概念。取而代之的是,流式应用需要针对到来数据 尽可能快的计算结果,同时还要应对 很高的事件接入速率,我们用 延迟吞吐 来表示这两方面的性能需求。

2.1.1 延迟

延迟表示处理一个事件所需的时间,本质上它是从接收事件到在输出中观察到事件处理效果的时间间隔。

在流处理中,延迟是以时间片(例如毫秒)为单位测量的。根据应用的不同,你可能会关注平均延迟,最大延迟或延迟的百分位数值。

保证低延迟,对很多流式应用(例如:诈骗识别、系统告警、网络监测,以及遵循服务级别协议的服务)而言至关重要。低延迟是流处理的一个关键特性,它滋生出了所谓的实时应用。

2.1.2 吞吐

吞吐是用来衡量系统处理能力(处理速率)的指标,它告诉我们系统 每单位时间可以处理多少事件

吞吐的衡量方式是计算每个单位时间的事件或操作数。但要注意,处理速率取决于数据到来速率,因此吞吐低不一定意味着性能差。在流处理系统中,你通常希望系统有能力应对以最大期望速率到来的事件。换言之,首要的关注点是确定 峰值吞吐,即系统满负载时的性能上限。

2.1.3 延迟与吞吐

延迟和吞吐并非相互独立的指标。如果事件在数据处理管道中传输时间太久,我们将难以确保高吞吐;同样,如果系统性能不足,事件很容易堆积缓冲,必须等待一段时间才能处理。然而,通过并行处理多条数据流,可以在处理更多事件的同时降低延迟。

2.2 数据流上的操作

流处理引擎通常会提供一系列内置操作来实现数据流的获取、转换,以及输出。这些算子可以组合生成 Dataflow 处理图,从而实现流式应用所需的逻辑。

这些操作既可以是 无状态stateless)的,也可以是 有状态stateful)的。无状态的操作不会维持内部状态,即处理事件时无需依赖已处理过的事件,也不保存历史数据。由于事件处理互不影响且与事件带来的时间无关,无状态的操作很容易并行化。此外,如果发生故障,无状态的算子可以很容易地重启,并从中断处继续工作。相反,有状态算子可能需要维护之前接收的事件信息。它们的状态会根据传入的事件更新,并用于未来事件的处理逻辑中。有状态的流处理应用在并行化和容错方面会更具挑战性,因为它们需要对状态进行高效划分,并且在出错时需进行可靠的故障恢复。

2.2.1 数据接入和数据输出

数据接入和数据输出操作允许流处理引擎和外部系统进行通信。数据接入操作是从外部数据源获取原始数据并将其转换成适合后续处理的格式。实现数据接入操作逻辑的算子称为 数据源。数据源可以从 TCP 套接字、文件、Kafka 主题或传感器数据接口中获取数据。数据输出操作是将数据嗯以适合外部系统使用的格式输出。负责数据输出的算子称为 数据汇,其写入的目标可以是文件、数据库、消息队列或监控接口等。

2.2.2 转换操作

转换操作是一类 只过一次 的操作,它们会分别处理每个事件。这些操作逐个读取事件,对其应用某些转换并产生一条新的输出流。转换逻辑可以是算子内置的,也可以由用户自定义函数提供。
在这里插入图片描述
算子既可以同时接收多个输入流或产生多条输出流,也可以通过单流分割或合并多条流来改变 Dataflow 图的结构。

2.2.3 滚动聚合

滚动聚合(如求和、求最小值和求最大值)会根据每个到来的事件持续更新结果。聚合操作都是有状态的,它们通过将新到来的事件合并到已有状态来生成更新后的聚合值。注意,为了更有效的合并事件和当前状态并生成单个结果,聚合函数必须满足 可结合可交换 的条件,否则算子就要存储整个流的历史记录。下图展示了一个求最小值的滚动聚合,其算子会维护当前的最小值,并根据每个到来的事件去更新这个值。
在这里插入图片描述

2.2.4 窗口操作

转换操作和滚动聚合每次处理一个事件来产生输出并(可能)更新状态。然而,有些操作必须收集并缓冲记录才能计算结果。例如流式 Join 或像是求中位数的整体聚合(holistic aggregate)。为了在无限数据流上高效地执行这些操作,必须对操作所维持的数据量加以限制。窗口操作 支持这项功能。

除了产生单个有用的结果,窗口操作还支持在数据流上完成一些具有切实语义价值的查询。你已经了解滚动聚合是如何将整条历史流压缩成一个聚合值,以及如何针对每个事件在极低延迟内产生结果。该操作对某些应用而言是可行的,但如果你只对最新的那部分数据感兴趣该怎么办呢?

窗口操作会持续创建一些称为 有限事件集合,并允许我们基于这些有限集进行计算。事件通常会根据其时间或其他数据属性分配到不同桶中。为了准确定义窗口算子语义,我们需要决定事件如何分配到桶中以及窗口用怎样的频率产生结果。窗口的行为是由一系列策略定义的,这些窗口策略决定了 什么时间创建桶事件如何分配到桶中 以及 桶内数据什么时间参与计算

其中参与计算的决策会根据触发条件判定,当触发条件满足时,桶内数据会发送给一个 计算函数evolution function),由它来对桶中的元素应用计算逻辑。这些计算函数可以是某些聚合(例如求和,求最小值),也可以是一些直接作用于桶内收集元素的自定义操作。策略的指定可以基于时间(例如最近 5 秒钟接收的事件)、数量(例如最新 100 个事件)或其他数据属性。

2.2.4.1 滚动窗口

滚动窗口(tumbling window)将事件分配到长度固定且互不重叠的桶中,在窗口边界通过后,所有事件会发送给计算函数进行处理。

基于数量的滚动窗口 定义了在触发计算器需要集齐多少条事件。

在这里插入图片描述
基于时间的滚动窗口 定义了在桶中缓冲数据的时间间隔。
在这里插入图片描述

2.2.4.2 滑动窗口

滑动窗口(sliding window)将事件分配到大小固定且允许相互重叠的桶中,这意味着每个事件可能会同时属于多个桶。我们通过指定长度(fixed length)和滑动间隔(slide)来定义滑动窗口。滑动间隔决定每隔多久生成一个新的桶。

在这里插入图片描述
上图为长度为 4 个事件、滑动间隔为 3 个事件的基于数量的滑动窗口。

2.2.4.3 会话窗口

会话窗口(session window)在一些常见的真实场景中非常有用,这些场景既不适合用滚动窗口,也不适合用滑动窗口。假设有一个应用要在线分析用户行为,在该应用中我们要把事件按照用户的同一活动或会话来源进行分组。会话由发生在相邻时间内的一系列事件,外加一段非活动时间组成。例如,用户浏览一连串新闻文章的交互过程,可以看做一个会话。由于会话长度并非预先定义好,而是和实际数据有关,所以无论是滚动还是滑动窗口都无法用于该场景。而我们需要一个窗口操作,能将属于同一会话事件分配到相同桶中。会话窗口根据会话间隔(session gap)将事件分为不同的会话,该间隔值定义了绘画在关闭前的非活动时间长度。

在这里插入图片描述

2.2.4.4 小结

迄今为止,你所见到的所有窗口都是基于 全局流数据 的窗口。但在实际应用中,你可能会想将数据流划分为多条逻辑流并定义一些并行窗口。例如,如果你在收集来自不同传感器的测量值,那么可能会想在应用窗口计算器按照传感器 ID 对数据流进行划分。并行窗口中,每个数据分区所应用的窗口策略都相互独立。下图展示了一个按事件颜色划分、基于数量 2 的并行滚动窗口。

在这里插入图片描述

窗口操作与流处理中两个核心概念密切相关:时间语义time semantics)和 状态管理state management)。时间可能是流处理中最重要的一个方面。尽管低延迟是流处理中一个很吸引人的特性,但流处理的真正价值远不止提供快速分析。

现实世界的系统、网络及通信信道往往充斥着缺陷,因此流数据通常都会有所延迟或者以乱序到达。了解如何在这种情况下提供精准确定的结果就变得至关重要。此外,处理实时事件的流处理应用还应以相同的方式处理历史事件,这样才能支持离线分析,甚至时间旅行式分析(time travel analysis)。当然,如果你的系统无法在故障时保护状态,那一切都是空谈。

至今为止你见到的所有窗口类型都要在生成结果前缓冲数据。实际上,如果你想在流式应用中计算任何有意义的结果(即便是简单的计数),都需要维护状态。考虑到流式应用可能需要整日、甚至长年累月的运行,因此必须保证出错时其状态能进行可靠的恢复,并且即使系统发生故障,系统也能提供准确的结果。后续,我们将深入研究流处理中的时间以及发生故障时和状态保障相关的概念。

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

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

相关文章

【江科大】STM32:(超级详细)定时器输出比较

文章目录 输出比较单元特点 高级定时器:均有4个通道 PWM简介PWM(Pulse Width Modulation)脉冲宽度调制输出比较通道PWM基本结构基本定时器 参数计算捕获/比较通道的输出部分详细介绍如下: 舵机介绍硬件电路 直流电机介绍&#xff…

LLM自回归解码

在自然语言处理(NLP)中,大型语言模型(LLM)如Transformer进行推理时,自回归解码是一种生成文本的方式。在自回归解码中,模型在生成下一个单词时会依赖于它之前生成的单词。 使用自回归解码的公式…

SPE-Single Pair Ethernet单对以太网测试那些事儿

SPE-Single Pair Ethernet单对以太网测试哪些事?SPE标准IEEE802.3再网上溯源的话是从ISO/IEC11801-X series演变而来。 IEEE802.3cg 10Base-T1 10mbt/s 15m-1000m 0.1mHz-20mHz IEEE802.3bw 100Base-T1 100mbt/s 15m 0.3mHz-66mHz IEEE802.3bp 1000…

k8s-认证授权 14

Kubernetes的认证授权分为认证(鉴定用户身份)、授权(操作权限许可鉴别)、准入控制(资源对象操作时实现更精细的许可检查)三个阶段。 Authentication(认证) 认证方式现共有8种&…

Pandas.Series.describe() 统计学描述 详解 含代码 含测试数据集 随Pandas版本持续更新

关于Pandas版本: 本文基于 pandas2.1.2 编写。 关于本文内容更新: 随着pandas的stable版本更迭,本文持续更新,不断完善补充。 传送门: Pandas API参考目录 传送门: Pandas 版本更新及新特性 传送门&…

Java层序遍历二叉树

二叉树准备: public class TreeNode {int val;TreeNode left;TreeNode right;TreeNode() {}TreeNode(int val) {this.val val;}TreeNode(int val, TreeNode left, TreeNode right) {this.val val;this.left left;this.right right;} } 思路:我们需要创建一个队…

前后端分离,使用vue3整合SpringSecurity加JWT实现登录校验

前段时间写了一篇spring security的详细入门,但是没有联系实际。 所以这次在真实的项目中来演示一下怎样使用springsecurity来实现我们最常用的登录校验。本次演示使用现在市面上最常见的开发方式,前后端分离开发。前端使用vue3进行构建,用到…

算法每日一题: 分割数组的最大值 | 动归 | 分割数组 | 贪心+二分

Hello,大家好,我是星恒 呜呜呜,今天给大家带来的又是一道经典的动归难题。 题目:leetcode 410给定一个非负整数数组 nums 和一个整数 k ,你需要将这个数组分成 k_ 个非空的连续子数组。设计一个算法使得这 k _个子数组…

Mybatis 动态SQL(set)

我们先用XML的方式实现 : 把 id 为 13 的那一行的 username 改为 ip 创建一个接口 UserInfo2Mapper ,然后在接口中声明该方法 package com.example.mybatisdemo.mapper; import com.example.mybatisdemo.model.UserInfo; import org.apache.ibatis.annotations.*; import jav…

mybatis的缓存机制

视频教程_免费高速下载|百度网盘-分享无限制 (baidu.com) MyBatis 有一套灵活而强大的缓存机制,主要分为两级缓存:一级缓存(本地缓存)和二级缓存(全局缓存)。 一级缓存(本地缓存)&a…

【网络奇遇记】揭秘计算机网络性能指标:全面指南

🌈个人主页:聆风吟 🔥系列专栏:网络奇遇记、数据结构 🔖少年有梦不应止于心动,更要付诸行动。 文章目录 📋前言一. 速率1.1 数据量1.2 速率 二. 带宽三. 吞吐量四. 时延4.1 发送时延4.2 传播时延…

PCIe-6328 八口USB3.0图像采集卡:专为工业自动化和机器视觉设计

PCIe-6328一块8口USB 3.0主控卡,专为工业自动化和机器视觉相关应用设计。USB 3.0或称作高速USB,是一项新兴总线技术,10倍于USB2.0的传输速度,尤其适用于高速数据存储和图 像设备。 绝大多数现有USB 3.0卡兼用多个接口于一个USB 3…

5. 函数调用过程汇编分析

函数调用约定 __cdecl 调用方式 __stdcall 调用方式 __fastcall 调用方式 函数调用栈帧分析 补充说明 不同的编译器实现不一样,上述情况只是VC6.0的编译实现即便是在同一个编译器,开启优化和关闭优化也不一样即便是同一个编译器同一种模式,3…

光催化专用设备太阳光模拟器装置

什么是光催化材料? 光催化材料是指通过该材料、在光的作用下发生的光化学反应所需的一类半导体催化剂材料。半导体是一种介于导体和绝缘体之间的物质,它有一个特殊的能带结构,即价带和导带之间有一个禁带,禁带的宽度决定了半导体…

linux压缩包形式安装mysql5.7

1. 下载 MySQL 压缩包 在官方网站或者镜像站下载 MySQL 压缩包。mysql-5.7.29-linux-glibc2.12.tar 下载地址: MySQL :: Download MySQL Community Server (Archived Versions) 2. 解压缩文件 使用以下命令解压 MySQL 压缩包: tar xvf mysql-5.7.29…

软考高项论文范文 | 进度管理

2017年5月,受某政府部门的委托,我单位承接了某信息共享与服务系统的建设工作,在本项目中我担任项目经理,负责项目的整体规划、组织实施和管理控制。某政府部门拥有多年积累的大量工作资源,这些信息是其开展各项行业服务…

南南合作里程碑!批量苏州金龙纯电公交正式交付哥斯达黎加

1月17日,哥斯达黎加电力研究所(ICE)收到了中国援助的批量苏州金龙海格纯电公交,该批车是中国应对气候变化南南合作援助哥斯达黎加项目的重要物资之一。中国驻哥斯达黎加特命全权大使汤恒出席交付仪式。 中国驻哥斯达黎加大使汤恒&…

esp8266小车智能wifi小车寒假营实战背篼酥老师

esp8266小车智能wifi小车寒假营实战 10节课 整车效果图如下 第一课 esp8266开发环境搭建和库文件加载 课程如下: 环境搭建 库文件下载链接:见文章末尾 第二课 小车模块组成和例程简介 课程如下: 车身PCB 小车电机 esp8266扩展板 esp8…

Flink实战之DataStream API

接上文:Flink实战之运行架构 Flink的计算功能非常强大,提供的应用API也非常丰富。整体上来说,可以分为DataStreamAPI,DataSet API 和 Table与SQL API三大部分。 其中DataStream API是Flink中主要进行流计算的模块。 DateSet API是…

vue3-组件通信

1. 父传子-defineProps 父组件&#xff1a; <script setup> import { ref } from vueimport SonCom from /components/Son-com.vueconst message ref(parent)</script><template><div></div><SonCom car"沃尔沃" :message"…