版权说明:本文档由用户提供并上传,收益归属内容提供方,若内容存在侵权,请进行举报或认领
文档简介
大数据处理
2025春#5:大数据流处理1目录5.1流处理基础和应用5.2分布式流计算5.3开源系统及编程思想5.4流处理系统机制及优化2目录5.1流处理基础和应用5.2分布式流计算5.3开源系统及编程思想5.4流处理系统机制及优化35.1流处理基础流处理概述流数据:流数据是一个无序的数据序列。在现实生活中,流数据无处不在,人们平常生活相关的所有信息都是随着时间推移不断产生的,它们即构成不同流数据。流处理:流处理技术是一门用于对流数据进行分析处理的技术,它以一个或多个流数据为输入,通过计算机技术分析产生有用信息。算子算子算子算子数据源算子算子算子算子45.1流处理基础TweetSpoutParseTweetBoltWordCountBolt在大数据时代,由于大量的数据以极快的速率持续产生,流处理技术更多是指对流数据进行处理的并行编程范式,它采用流水线思想将处理逻辑分解为多个处理步骤,并依次对数据流中的每一个元素进行处理。Twitter中执行WordCount应用时采用了流处理技术,将WordCount的过程分为三个流水线步骤:TweetSpout、ParseTweetBolt、WordCountBoltTweetSpout:源源不断从Twitter的接口中读取推文,然后发送给ParseTweetBoltParseTweetBolt:将接收到的推文分解为一个个的单词,然后发送给WordCountBoltWordCountBolt:对接收到的单词进行计数55.1流处理基础为什么需要流处理技术流处理技术出现之前,传统的对流数据的处理方法一般是:从应用、传感器等终端获取流数据并先存储在数据库、文件系统等静态存储系统中,然后在这些静态数据上进行查询或计算,这种处理实质上是采用了批处理技术批处理对数据流的处理过程应用其他传感器事件存储数据库分布式文件系统查询、更新查找分析数据以事件流的形式产生空闲时存储数据应用对数据进行计算65.1流处理基础为什么需要流处理技术流处理技术可以实现对数据流中的每一个元素进行实时处理如右图对比批处理过程,流处理省略了将数据存储在静态数据库或文件系统中的过程,在流数据不断产生的同时,应用也在不断地对数据进行查询、处理、分析等操作,并连续不断地产生结果。在这种架构下,应用处理的延时大大降低流处理系统ApacheStormSparkStreamingApacheFlink批处理对数据流的处理过程事件分布式文件系统数据以事件流的形式产生空闲时存储数据应用对数据进行计算应用传感器其他存储数据库查询、更新查找分析流处理对数据流的处理过程VS事件数据以事件流的形式产生应用传感器其他应用应用应用应用对数据进行计算75.1流处理基础流处理的需求在大数据时代,系统要处理的数据量非常大,以至于常规的处理模式无法在有效时间内处理完成。同时,数据产生速度快,例如微博每时每刻都有大量的新数据产生,由此产生数据流的速率非常快。大规模流处理系统主要需求:高吞吐低延迟可扩展高可用85.1流处理基础流处理应用:个性化推荐搜索引擎在生活中扮演着重要的角色,作为互联网上最大的流量入口,每天为无数条搜索请求返回结果。在人们使用搜索引擎进行搜索的同时,搜索引擎也得到了大量的用户数据。其中包括用户搜索历史信息,网页点击信息和位置信息等。通过这些信息,进行分析和对比,可以对用户的行为进行深入了解,推断用户的偏好信息,以提供个性化的服务,例如个性化的广告或新闻推荐例如广告商通过实时分析搜索引擎获取到的用户数据,有针对性地投放广告,提高广告点击率,获取利润如图,用户在查询框输入了一个查询请求,随后又在搜索结果页面点击了一个广告。这样,搜索引擎中会产生两条数据流:查询流和广告点击流服务器接收到这两条数据流,通过对相近时间内不同数据流中的数据进行连接操作:判断广告与此次查询相关性,提出相识查询用户可能也会点击此类广告搜索引擎中的流连接操作查询双击ID时间搜索引擎连接的事件95.1流处理基础流处理应用:实时在线交易上市公司通过在股票市场中发行股票来为公司筹集资本。股票只有在交易中体现价格。股票市场中进行的交易也是通过流处理来实现的,存在两条流:买入流和卖出流股票市场将两条流进行匹配,并进行交易,股票由卖出方转让到买入方,资金则相反。基于时间优先的竞价规则,流数据的顺序对结果影响尤为重要,同时,由于股票价格波动且涉及大量财产的交易,对于延迟和准确率的要求也比较高股票交易中的买入和卖出流买入价131517卖出价161918时间105.1流处理基础流处理应用:实时热点检测社交媒体中的热点话题反应了现实生活中人们对不同事件关注程度的差异。通过分析热点话题的时空变化,可以了解舆论走向,理解热点话题的传播状况。传统的基于批处理的热点话题检测模型比较简单,主要是将一段时间内用户发表的推文等信息聚集并把其中的字词提取出来,为转发评论量等赋予权重,通过聚类算法进行处理得出实时的话题热点,但是难以得出实时的结果流处理则可以通过实时检测系统更快地发现热门话题,以便做出响应Twitter的热点话题检测系统中,用户推文流不断产生并发送到服务器上。在服务端存在文本处理、频率分析等组件。首先推文流被分割成小的文本段,之后流过频率分析组件,产生不同词语的出现频率。最终经过统计计算和筛选,流出系统的就是最近一段时间内的热点话题Twitter热点话题检测系统数据获取Twitter热点话题输出文本分析筛选热点话题TwitterTwitter热点话题检测系统热点1热点2热点3......热点n115.1流处理基础流处理应用:欺诈检测随着互联网移动通信的普及,电信欺诈的例子屡见不鲜,犯罪分子通过伪装身份向用户传递伪造的信息,并通过引导,最终使得用户在不知不觉中将财产转移到犯罪分子的账户骗子A首先通过群发等方式向众多群众发出消息,告知其可能有某些行为违反了规定,并要求尽快同“相关群众”取得进一步的联系。这里的“相关人员”则是骗子A的同伙B。一部分群众因为没有清楚地辨别出骗子,就会听信骗子的话,然后再一步步地跟着骗子的引导,最终将钱转入骗子的账户。欺诈检测的目的就是实时检测出这张欺诈事件,以避免广大群众遭受财产损失。由于“广撒网”策略的存在,骗子A会通过同一个号码向多个号码发信息,而受骗者通常会打给相同的骗子B。在电信公司的通话数据流中对欺诈的特定模式进行匹配,实现识别。电信欺诈示意图骗子A骗子B1:M(M>>N)(N:1)群众12目录5.1流处理基础和应用5.2分布式流计算5.3开源系统及编程思想5.4流处理系统机制及优化135.2分布式流计算分布式流计算将流应用部署在分布式集群中,充分利用分布式集群的高并行计算能力,对连续到来的数据流进行不间断计算的过程。核心是将流应用逻辑构建成这样一个具有固定的起始、终止和若干中间操作的流水线计算模型,以数据流的形式传入系统当中。中间操作根据应用逻辑,将流中的数据进行过滤、计数、求和、连接等多种计算和转化,并将数据流不断向下一个操作传输,直到进行完终止操作,将应用的最终结果进行输出145.2分布式流计算分布式流计算重要步骤数据封装,将各种来源的待处理数据封装为连续的数据流模式建立应用拓扑,将应用逻辑转化为一组由起始、终止和一系列中间操作组成的应用拓扑图指定操作的并行度,为了充分利用分布式集群的高并行计算能力,每个操作在系统中可以由多个线程来同时运行,提高计算的吞吐率指定数据的分组和传输方式,由于操作多实例化,需要进一步指定流数据的分组和传输方式,以保证所有数据能够被正确和完整的处理155.2分布式流计算数据封装进行流处理,首先需要从外部数据源(如文件系统、时序数据库、监听网络socket等)获取数据并转化封装为统一的可以被流处理系统所接收的格式。需要将数据切分和转化为可以为一个个连续的可以被独立处理的最小数据单元,一个最小的可以被独立处理的数据单元被称为元组(tuple)每个元组表示为若干值的列表t=<value1,value2,value3,…>,每一项代表对应元组数据的一个固定属性值,列表长度由应用需求来确定例股票交易,一个元组包含交易时间、股票ID、交易金额、交易数目、交易对象等对个属性值165.2分布式流计算建立应用拓扑将复杂的应用逻辑拆分成若干简单而独立的计算操作。转化为一组由起始、终止和一系列中间操作组成的应用拓扑图。根据划分的操作和操作间的关联关系,流处理逻辑被建模成一个有向无环图的形式Foreachtuplet0=A(DataStream){t1=B(t0);t2=C(t0);If(condition1){t3=D(t1);Returnt3;}Else{t4=E(t1,t2);Returnt4;}}操作A操作B操作C操作D操作E此计算模型使得分布式流处理系统能够高效并行地对数据进行处理一个示例的流应用逻辑算法175.2分布式流计算指定操作的并行度在划分操作,建立完流应用拓扑后,可进一步指定每个操作的并行度。即在分布式系统中,对于每一个操作,都可以生成多个实例,对应多个线程来执行这一操作。操作对应的线程数越多,那么执行的并行度也越高,在这一操作中,数据流的平均处理速度也越快流处理中的操作多实例化流水线并行处理方式元组a元组b元组c元组a元组a元组b元组c元组a元组b操作A操作B操作D时间操作A操作E操作D操作BC1C2C3185.2分布式流计算指定数据分组与传输方式由于分布式流处理系统中的操作具有多实例化,所以上下游操作间的一对一关系演变成实例与实例的多对多关系。为了保证数据在不同操作间传输时不会被重复处理或者产生遗漏,还需要指定有效的数据分组策略:Shufflegrouping将数据以轮廓或随机的方式进行分组,即无视数据本身值的信息,将数据划分为相近规模的多个子数据流Keygrouping将元组中的某一个属性值指定为其键,所有的元组都根据这个键值进行哈希映射,使得具有相同键值的数据可以被划分到下游操作中同一实例进行操作Allgrouping将数据进行广播发送,将上游产生的每个流数据元组复制和分发到下游每一个操作的每一个实例当中Globalgrouping将所有元组分组到同一个实例中19目录5.1流处理基础和应用5.2分布式流计算5.3开源系统及编程思想5.4流处理系统机制及优化205.3分布式系统及编程模型基于流计算的基本模型,当前已有各式各样的分布式流处理系统被开发出来:ApacheStormSparkStreamingApacheFlink215.3分布式系统及编程模型ApacheStorm由Twitter公司开源的一个实时分布式流处理系统,被广泛应用在实时分析、在线机器学习、连续计算、分布式RPC、ETL等场景支持水平扩展、高容错性,保证数据能被处理,而且处理速度很快支持多编程语言,易于部署225.3分布式系统及编程模型Storm:数据封装从分布式文件系统或分布式消息队列中获取源数据,并将每个流数据元组封装称为tuple。一条数据流即是一个无边界的tuple序列,而这些tuple序列可以以分布式的方式创建和处理。一个tuple可以包含多个字段,每个字段代表对于流数据的一个属性,在Storm的每个操作组件发送向下游发送tuple时,会声明对应tuple每个字段的顺序和代表的含义field1field2field3field4235.3分布式系统及编程模型Storm:应用拓扑建立用户所提交的应用所构建的DAG拓扑被称为Topology。Storm的Topology类似于MapReduce中的一个job,但区别在于这个拓扑会永远运行(或者直到手动结束)。每个Topology中有两个重要组件:spout和boltSpout是topology中数据流的来源,也即对应DAG模型中的起始操作。Spout可以从外部源读取数据并将其以封装成tuple的形式发送到Topology中Bolt是Topology中对tuple进行处理的主要单元。Storm并不区分中间和终止操作,而是将其统一为bolt来进行实现,也即对结果的输出需要由用户自己来实现245.3分布式系统及编程模型Storm:并行度指定Storm中并行度有三层含义:Storm可以建立在分布式集群上,每台物理节点可以发起一个或多个worker进程,一个worker对应一个物理的JVM整个Topology会由一个或者多个worker进程来负责执行,每个worker会在一个JVM中运行一个或多个executor,每个executor对应一个线程,执行某一个spout或者bolt的计算任务每个spout/blot都可以实例化生成多个task在集群中运行,一般默认情况下,executor数与task数一一对应,即每个实例都由一个单独的线程来执行255.3分布式系统及编程模型Storm:数据分组和传输用户通过定义分组策略(streaminggrouping)来决定数据流如何在不同的spout/bolt的task中进行分发和传输。分组策略将所有的spout和bolt连接起来构成一个Topology,如右图Localgrouping策略,是shufflegrouping的一种变种分组策略。由于Storm划分为多个worker进程,shufflegrouping可能导致大量的进程间通信,localgrouping则是将元组优先发往与自己同进程的下游task中,若没有这种下游task,才继续沿用shufflegrouping的方式directgrouping策略,是一种特殊的分组方式,用户可以直接指定由下游的哪一个task来接收数据streaminggroupingspoutboltAboltBboltC265.3分布式系统及编程模型Storm:分布式系统架构Storm可以运行在分布集群上。Storm集群结构沿用了主从架构方式,即一个主控节点和多个工作节点Storm的基本组件:Nimbus:主控节点Supervisor:接收Nimbus分配的任务ZooKeeper:进行Nimbus和Supervisor之间协调工作Storm系统架构NimbusZookeeperZookeeperZookeeperSupervisorSupervisorSupervisorSupervisorSupervisorworkersworkersworkersworkersworkers275.3分布式系统及编程模型Storm:WordCount应用编程示例实现生成数据的spout,封装数据首先构建一个CreateSentenceSpout来进行数据流的生成实现对流数据进行操作处理的bolt在WordCount应用中,对spout生成的句子,构建两个bolt来进行处理:一个SplitWordBolt来将句子划分为单词,一个CountBolt来对划分好的单词进行累计计数构建流应用Topology,并指明并行度和分组策略实现了对应的spout和bolt功能之后,最后就是将其连接成一个完整的Topology285.3分布式系统及编程模型SparkStreaming是SparkAPI核心扩展,提供对实时数据流进行流式处理,具备可扩展性、高吞吐和容错等特性。支持从多种数据源中提取数据,例如Twitter、Kafka、Flume、ZeroMQ和TCP套接字,并提供了高级API来表示复杂处理算法,如map、reduce、join等295.3分布式系统及编程模型SparkStreaming:数据封装本质上是一个典型的微处理系统,其与以元组为单位进行流式处理不同,他将无尽的数据流按时间切分为连续的小批次数据,然后以传统的批处理方法来进行快速连续的处理在SparkStreaming中,数据流被抽象成以时间片段分隔开的离散流形式SparkStreaming使用Spark引擎,将每一段小批次数据转化成为Spark当中的RDD。数据流以RDD形式在系统中进行运算SparkStreaming的离散流时间0到1的数据时间1到2的数据时间2到3的数据时间3到4的数据01234时间离散流RDDRDDRDDRDD305.3分布式系统及编程模型SparkStreaming:应用拓扑建立在系统中构建出DAG的处理模型。与Storm不同,SparkStreaming并不使用固定的处理单元来执行单一的操作,其DAG与SparkCore中的DAG相同,只是用DAG的形式将每一个时间分片对应的RDD进行运算的job来进一步划分成任务集stage沿用了SparkCore对RDD提供的transformation操作,将所有RDD依次进行转换,应用逻辑分别进行转换处理,进而实现对整个离散流的转换在线输入的数据流被按照时间切分为若干小批次数据并转化成为RDD存储在内存中;根据流应用逻辑,也即流处理引用抽象出DAG拓扑,指定出相应的RDDtransformationSparkStreaming计算框架将数据流划分为微批用离散流形式来表示流计算任务调度器内存管理器SparkSpark用批量job来执行RDDtransformations将每批输入数据转化为RDD产生RDDtransformationsSparkStreaming在线输入的数据流batchesofresults315.3分布式系统及编程模型SparkStreaming:并行度指定由于本质上是将数据流的任务划分成大量的微批数据,对应多个job来执行,所以其并行度设定与Spark进行批处理时的设定一样,只能设定整体job的并行度,而不能对每个操作独立的并行度进行设置。由于批处理的特性,SparkStreaming可以最大化对系统并行能力的利用,也能获得相对更高的系统吞吐率325.3分布式系统及编程模型SparkStreaming:数据分组和传输数据被打包为一个个微批,而每个微批相互独立地进行处理,所以不涉及所提到的数据分组与传输问题。但这也展现出微批处理的一个局限性,其难以灵活处理基于用户自定义的窗口的聚合、计数等操作,也不能进行针对数据流的连续计算,如两个数据流的实时连接等操作335.3分布式系统及编程模型SparkStreaming:系统框架SparkStreaming建立在Spark系统之上,其系统架构相对于Spark的修改和新增如图:SparkStreaming主要组件:Master:流应用入口Client:将数据传入到系统当中Worker:流数据的入口以及执行RDD转换的主要组件SparkStreaming基于Spark修改和新增的组件masterD-StreamlineageinputtrackerRDDlineagetaskschedulerblocktrackerinputreceivertaskexecutionblockmanagercomm.managerworkerinputreceivertaskexecutionblockmanagercomm.managerworkerclientclient输入及checkpoint的RDD的副本新增组件修改组件345.3分布式系统及编程模型SparkStreaming:编程示例离散流的输入和数据封装建立应用拓扑,进行离散流的转化355.3分布式系统及编程模型ApacheFlink一个同时支持分布式数据流处理和数据批处理的大数据处理系统完全以流处理的角度出发进行设计,而将批处理看作是有边界的流处理特殊流处理来执行Flink使用单纯流处理方法的典型系统,其计算框架与原理和storm比较相似Flink可以表达和执行许多类别的数据处理应用程序,包括实时数据分析、连续数据管道、历史数据处理和迭代算法365.3分布式系统及编程模型ApacheFlink:数据封装能够支撑对多种类型的数据进行处理类似Storm,Flink同样也可以使用多字段的tuple为其基本数据单元Flink可以支持多种Flinktuple类型,每种tuple都是一个固定长度的对象序列375.3分布式系统及编程模型ApacheFlink:应用拓扑建立Flink中核心概念为数据流和转换。每个转换对应的是一个简单的操作,根据应用逻辑,转换按先后顺序构成了流应用的DAG图。如图,数据流在转换之间传递,直到完成所有的转换进行输出ApacheFlink的计算模型sourcemap()KeyBy()/window()/apply()sink源操作转换操作汇聚操作StreamStreamingdataflow385.3分布式系统及编程模型ApacheFlink:并行度指定与Storm相似,Flink程序的计算框架本质上也是并行分布的。在系统中,一个流包含一个或多个流分区,而每一个转换操作包含一个或多个子任务实例。操作的子任务间彼此独立,以不同的线程执行,可以运行在不同的机器或容器上。395.3分布式系统及编程模型ApacheFlink:数据分组与传输包括一对一(one-to-one)模式或者重分组(redistributing)模式一对一模式,数据流中元素的分组和顺序会保持不变,也就是说,对于上下游的两个不同的转换操作,下游任一子任务内要处理的元数据,与上游相同顺序的子任务所处理的元组数据完全一致重分组模式,会改变数据流所在的分组。重分组后元组的目标子任务根据处理的变换方法不同而发生改变405.3分布式系统及编程模型ApacheFlink:系统框架jobclient:独立的程序执行入口jobmanager:对应一个Flink程序的master进程,负责job的管理和资源的协调taksmanager和taskslot:具体负责执行task的组件。每个taskmanager对应是运行在节点上的JVM进程ApacheFlink分布式运行环境程序代码optimizer/graphbuilderclientactorsystem程序dataflowdataflowgraphFlink程序jobmanageractorsystemdataflowgraph调度器检查点协调器taskmanagertaskslottasktaskslottasktaskslot内存&I/O管理器网络管理器actorsystemtaskmanagertaskslottasktaskslottasktaskslot内存&I/O管理器网络管理器actorsystem数据流(worker)(worker)(master/YARNapplicationmaster)Task状态心跳统计数据部署/停止/取消task触发检查点状态更新统计数据&结果提交job(senddataflow)取消/更新job415.4流处理系统机制及优化流处理调度及优化在分布式流处理系统中,为了保证流处理系统的高吞吐,不让任何一个处理单元成为系统处理的瓶颈,流处理系统需要合理的调度来保证各个处理单元的负载均衡。流处理调度通常包含任务和数据的调度:任务调度:将流处理应用产生的实例部署到集群的服务器过程数据调度:上下流操作的实例间的数据分发的方式425.4流处理系统机制及优化流处理任务调度任务调度通常可以定义为应用产生的实例集合到集群服务器的一种映射关系
π:
T→N,其中T={t1,t2,…}是指应用的所有计算实例构成的集合,N={n1,n2,…}是指集群中所有服务器构成的集合。在分布式流处理系统中,任务调度对系统的吞吐和处理延迟这两方面的性能起到了至关重要的作用。常见的一种任务调度策略是轮询(round-robin)调度,它将随机排好序的实例以轮调的方式依次部署到排好序的服务器上。Storm就以轮询调度作为默认的任务调度策略。435.4流处理系统机制及优化流处理数据调度对数据向下游多个实例进行调度时通常使用shufflegrouping和keygrouping两种策略,尤其对于本身具有键值的流数据,常常采用后者。然而keygrouping下数据的划分往往难以达到均衡,尤其是在数据具有高倾斜特征的情况下例如,微博流、股票交易流,大量的流数据都具有相同的键值。当分布式流处理系统在对这些流数据进行调度时,就会出现大量具有相同键值的数据被分组到同一个实例上的情况;相比之下,其他实例上的工作负载相比这些处理热门数据实例少很多,因而给系统带了严重的负载不平衡。流数据的keygrouping策略sourceWordCountWordCount...445.4流处理系统机制及优化流处理数据调度在面向高倾斜分布数据进行处理时,使用无关数据键值的轮询调度策略是一种基本的转换思想,即将流数据按到来的顺序依次以轮询的方式发送给各个实例来保证的完全负载均衡。轮询调度失去了对相同键值的数据处理的本地性。具有相同键值的数据会被分散在任一可能的实例上来进行处理,每一实例都会因此带来更高的内存和处理开销在键值调度策略下,增加实例个数可以使得每个实例平均需要处理的键值数量降低,从而减少内存开销。然而轮询调度则无法做到这一点,从而减少内存开销。当键值的总量到达一定规模时,甚至会使得系统出现存储资源无法满足应用需求的情况流数据的shufflegrouping策略sourceWordCountWordCount...聚合...455.4流处理系统机制及优化高倾斜数据流解决思路使用键值调度的思想,根据数据项的分布情况,重新设计键值空间划分的哈希函数,使得每个实例上都存在若干高频出现的数据项,降低系统实例间不均衡程度。使用多个哈希函数,使得每个键值都可以被映射到不同的实例之上,然后按照这些实例的当前负载情况,将数据进行合理的调度,避免集中到同一个实例之上实时的识别当前流数据中高频出现的若干数据项,对它们对应的数据采用轮询调度策略;对于其他数据则依然采用键值调度策略。基于此思想,Dstream是一个开源的基于热度感知分布式流处理区分调度系统Dstream系统架构概率计数器...概率计数衰减器潜在热词摘要热词预热器用户处理逻辑热词过滤器区分调度器......465.4流处理系统机制及优化流处理一致性语义流处理业务对数据需求多样:一方面,基于DAG计算模型会带来频繁的数据传输,同时流处理系统基于内存持续不断的计算,数据在传输和处理过程中存在丢失或失效的可能另一方面,部分流处理应用更注重处理的时效性,而并不在意是否有数据丢失,而部分流处理应用必须保证每个流数据元组都不能被丢弃流处理系统需要提供各种一致性语义来满足应用需求475.4流处理系统机制及优化流处理一致性语义流处理系统的语义保障通常按照每一个元组被完全处理的次数分为:至多一次(atmostonce)至少一次(atleastonce)恰好一次(exactlyonce)一个元组被完全处理指的是该元组以及由元组经过DAG计算模型中对应的所有操作处理后生成的中间元组均已被处理结束,如下图元组处理过程元组源操作中间处理操作中间处理操作汇聚操作485.4流处理系统机制及优化三种处理语义至多一次:指系统能够保障在处理数据时,任意元组要么刚好被处理一次,要么被丢弃,不会出现对同一元组的重复处理适用于统计分析关联的应用,对数据处理的完整性要求不高至少一次:指系统能够保障处理数据时,任意元组都至少被处理一次不能丢失,但是允许同一元组被重复处理适用于安全监测等相关应用,更要求处理的完整性和安全性,要求数据不能产生遗漏,即使产生误判也不可发生漏判恰好一次:指系统能够保障数据处理时,所有元组都恰好被处理一次,不存在元组被重复处理或者被丢弃的情况适用于金融交易等应用中,对数据的处理非常敏感。除了对数据进行持久化存储外,也需要对数据处理过程中产生的所有状态信息进行持久化存储495.4流处理系统机制及优化Storm采用的轻量级的ACK机制来实现至少一次语义当每一个元组进入到Storm时,它的数据会被持久化保存,同时acker会持续追踪这个元组,并为它记录一个校验值。
对于某个元组t1,当其进入到分布式流处理系统中时,acker记录一个随机生成的校验值α。当其在源操作中完成处理,并转化成为一个新元组t2,且
温馨提示
- 1. 本站所有资源如无特殊说明,都需要本地电脑安装OFFICE2007和PDF阅读器。图纸软件为CAD,CAXA,PROE,UG,SolidWorks等.压缩文件请下载最新的WinRAR软件解压。
- 2. 本站的文档不包含任何第三方提供的附件图纸等,如果需要附件,请联系上传者。文件的所有权益归上传用户所有。
- 3. 本站RAR压缩包中若带图纸,网页内容里面会有图纸预览,若没有图纸预览就没有图纸。
- 4. 未经权益所有人同意不得将文件中的内容挪作商业或盈利用途。
- 5. 人人文库网仅提供信息存储空间,仅对用户上传内容的表现方式做保护处理,对用户上传分享的文档内容本身不做任何修改或编辑,并不能对任何下载内容负责。
- 6. 下载文件中如有侵权或不适当内容,请与我们联系,我们立即纠正。
- 7. 本站不保证下载资源的准确性、安全性和完整性, 同时也不承担用户因使用这些下载资源对自己和他人造成任何形式的伤害或损失。
最新文档
- 2025绍兴诸壁中学高一数学分班考试真题含答案
- 2026年广东专升本数学考试题库及参考答案
- 2026年贵州公务员行测真题试卷带答案
- 2025~2026学年黑龙江哈尔滨市第六十九中学统编版七年级下学期学情检测历史试卷
- 2025~2026学年辽宁省锦州市第四中学教育集团七年级下学期期中测试历史试卷
- 2026年高中语文《蜀相》杜甫凭吊先贤咏史教案
- 2026年高中语文《桂枝香·金陵怀古》登高怀古教案
- 2026年坚壁清野成语故事历史典故拓展教案
- 2026年手不释卷成语故事阅读习惯教案
- 保健学试题与准确答案
- 企业精细化管理实施方案
- 小学数学人教版(新教材)五年级上还原简单组合体课件(共26张)
- 妊娠期高血压疾病诊治指南解读 课件
- 2026广东广州市南沙区黄阁镇人民政府招聘编外工作人员10人考前冲刺试卷附答案详解(研优卷)
- 2026年陕西省延安市重点学校初一入学数学分班考试试题及答案
- 超市连锁2026年员工劳动合同模板
- 2026年信息处理技术员(基础知识、应用技术)合卷软件资格考试(初级)试题附答案
- 立法研究基地工作方案
- 2026年秋人教PEP版(新教材)小学英语六年级上册《Unit 6 Energy,nature and us》单元达标自测卷及答案
- 2026年中职焊接(电阻焊)试题及答案
- 2026人教版四年级数学上册第一单元第2课《亿以内数的读法》课件
评论
0/150
提交评论