发布时间:2026/8/27 21:09:47
头歌实践教学平台:数据科学与大数据技术导论(十一下) 十一、数据处理与分析第3关大数据处理与分析技术任务描述本关任务根据相关知识内容完成右边选择题。相关知识为了完成本关任务你需要掌握1.技术分类2.流计算3.图计算。技术分类大数据计算模式及其代表产品大数据计算模式 解决问题 代表产品批处理计算 针对大规模数据的批量处理 MapReduce、Spark等流计算 针对流数据的实时计算 Flink、Storm、S4、SparkStreaming、Flume、Streams、Puma、DStream、Super Mario、银河流数据处理平台等图计算 针对大规模图结构数据的处理 Pregel、GraphX、Giraph、PowerGraph、Hama、GoldenOrb等查询分析计算 大规模数据的存储管理和查询分析 Dremel、Hive、Cassandra、Impala等流计算流计算是实时获取来自不同数据源的海量数据经过实时分析处理获得有价值的信息。流计算秉承一个基本理念即数据的价值随着时间的流逝而降低如用户点击流。因此当事件出现时就应该立即进行处理而不是缓存起来进行批量处理。为了及时处理流数据就需要一个低延迟、可扩展、高可靠的处理引擎对于一个流计算系统来说它应达到如下需求1高性能处理大数据的基本要求如每秒处理几十万条数据。2海量式支持 TB 级甚至是 PB 级的数据规模。3实时性保证较低的延迟时间达到秒级别甚至是毫秒级别。4分布式支持大数据的基本架构必须能够平滑扩展。5易用性能够快速进行开发和部署。6可靠性能可靠地处理流数据。传统的数据处理流程需要先采集数据并存储在关系数据库等数据管理系统中之后由用户通过查询操作和数据管理系统进行交互。传统的数据处理流程隐含了两个前提1存储的数据是旧的。存储的静态数据是过去某一时刻的快照这些数据在查询时可能已不具备时效性了2需要用户主动发出查询来获取结果流计算的处理流程流计算的处理流程一般包含三个阶段数据实时采集、数据实时计算、实时查询服务。1、数据实时采集数据实时采集阶段通常采集多个数据源的海量数据需要保证实时性、低延迟与稳定可靠。以日志数据为例由于分布式集群的广泛应用数据分散存储在不同的机器上因此需要实时汇总来自不同机器上的日志数据。目前有许多互联网公司发布的开源分布式日志采集系统均可满足每秒数百MB的数据采集和传输需求如Facebook 的 ScribeLinkedIn 的 Kafka淘宝的 Time Tunnel基于 Hadoop 的 Chukwa 和 Flume。数据采集系统的基本架构一般有以下三个部分1Agent主动采集数据并把数据推送到 Collector 部分。2Collector接收多个 Agent 的数据并实现有序、可靠、高性能的转发。3Store存储 Collector 转发过来的数据对于流计算不存储数据。如下图所示数据采集系统基本架构2、数据实时计算数据实时计算阶段对采集的数据进行实时的分析和计算并反馈实时结果。经流处理系统处理后的数据可视情况进行存储以便之后再进行分析计算。在时效性要求较高的场景中处理之后的数据也可以直接丢弃。3、实时查询服务经由流计算框架得出的结果可供用户进行实时查询、展示或储存。传统的数据处理流程用户需要主动发出查询才能获得想要的结果。而在流处理流程中实时查询服务可以不断更新结果并将用户所需的结果实时推送给用户。虽然通过对传统的数据处理系统进行定时查询也可以实现不断地更新结果和结果推送但通过这样的方式获取的结果仍然是根据过去某一时刻的数据得到的结果与实时结果有着本质的区别。可见流处理系统与传统的数据处理系统有如下不同流处理系统处理的是实时的数据而传统的数据处理系统处理的是预先存储好的静态数据。用户通过流处理系统获取的是实时结果而通过传统的数据处理系统获取的是过去某一时刻的结果。流处理系统无需用户主动发出查询实时查询服务可以主动将实时结果推送给用户。图计算传统图计算解决方案的不足之处。很多传统的图计算算法都存在以下几个典型问题1常常表现出比较差的内存访问局部性。2针对单个顶点的处理工作过少。3计算过程中伴随着并行度的改变。针对大型图比如社交网络和网络图的计算问题可能的解决方案及其不足之处具体如下1为特定的图应用定制相应的分布式实现通用性不好。2基于现有的分布式计算平台进行图计算在性能和易用性方面往往无法达到最优。现有的并行计算框架像 MapReduce 还无法满足复杂的关联性计算。MapReduce 作为单输入、两阶段、粗粒度数据并行的分布式计算框架在表达多迭代、稀疏结构和细粒度数据时力不从心。比如有公司利用 MapReduce 进行社交用户推荐对于 5000 万注册用户50 亿关系对利用 10 台机器的集群需要超过 10 个小时的计算。3使用单机的图算法库比如 BGL、LEAD、NetworkX、JDSL、Standford GraphBase 和 FGL 等但是在可以解决的问题的规模方面具有很大的局限性。4使用已有的并行图计算系统比如Parallel BGL 和 CGM Graph实现了很多并行图算法但是对大规模分布式系统非常重要的一些方面比如容错无法提供较好的支持。传统的图计算解决方案无法解决大型图的计算问题因此就需要设计能够用来解决这些问题的通用图计算软件。针对大型图的计算目前通用的图计算软件主要包括两种第一种主要是基于遍历算法的、实时的图数据库如 Neo4j、OrientDB、DEX 和 Infinite Graph。第二种则是以图顶点为中心的、基于消息传递批处理的并行引擎如 GoldenOrb、Giraph、Pregel 和 Hama这些图处理软件主要是基于 BSP 模型实现的并行图处理系统。一次 BSP (Bulk Synchronous Parallel Computing Model又称“大同步”模型)计算过程包括一系列全局超步所谓的超步就是计算中的一次迭代每个超步主要包括三个组件1局部计算每个参与的处理器都有自身的计算任务它们只读取存储在本地内存中的值不同处理器的计算任务都是异步并且独立的。2通讯处理器群相互交换数据交换的形式是由一方发起推送 (put) 和获取 (get) 操作。3栅栏同步Barrier Synchronization当一个处理器遇到“路障”或栅栏会等到其他所有处理器完成它们的计算步骤每一次同步也是一个超步的完成和下一个超步的开始 。编程要求根据相关知识完成右侧选择题测试说明若选择题答案与正确答案一致则可通关。开始你的任务吧祝你成功第4关大数据处理与分析代表性产品任务描述本关任务根据相关知识内容完成右边选择题。相关知识为了完成本关任务你需要掌握1.分布式计算框架 MapReduce2.数据仓库 Hive3.数据仓库 Impala4.基于内存的分布式计算框架 Spark5.TensorFlowOnSpark6.流计算框架 Storm7.流计算框架 Flink8.大数据编程框架 Beam9.查询分析系统 Dremel。分布式计算框架 MapReduceMapReduce 简介谷歌在 2003 年至 2006 年连续发表了 3 篇很有影响力的文章分别阐述了 GFS、MapReduce 和 BigTable 的核心思想。其中 MapReduce 是谷歌公司的核心计算模型。MapReduce 将复杂的、运行于大规模集群上的并行计算过程高度地抽象到两个函数Map 和 Reduce这两个函数及其核心思想都源自函数式编程语言。在 MapReduce 中一个存储在分布式文件系统中的大规模数据集会被切分成许多独立的小数据块这些小数据块可以被多个 Map 任务并行处理。MapReduce 框架会为每个 Map 任务输入一个数据子集Map 任务生成的结果会继续作为 Reduce 任务的输入最终由 Reduce 任务输出最后结果并写入分布式文件系统。特别需要注意的是适合用 MapReduce 来处理的数据集需要满足一个前提条件待处理的数据集可以分解成许多小的数据集而且每一个小数据集都可以完全并行地进行处理。MapReduce 设计的一个理念就是“计算向数据靠拢”而不是“数据向计算靠拢”因为移动数据需要大量的网络传输开销尤其是在大规模数据环境下这种开销尤为惊人所以移动计算要比移动数据更加经济。本着这个理念在一个集群中只要有可能MapReduce 框架就会将 Map 程序就近地在 HDFS 数据所在的节点运行即将计算节点和存储节点放在一起运行从而减少了节点间的数据移动开销。MapReduce 工作流程MapReduce 的工作流程中不同的 Map 任务之间不会进行通信不同的 Reduce 任务之间也不会发生任何信息交换用户不能显式地从一台机器向另一台机器发送消息所有的数据交换都是通过 MapReduce 框架自身去实现的。如下图所示MapReduce 工作流程MapReduce 的不足之处总体而言Hadoop 的 MapReduce 存在以下缺点。1表达能力有限计算都必须要转化成 Map 和 Reduce 两个操作但这并不适合所有的情况难以描述复杂的数据处理过程。2 磁盘 IO 开销大每次执行时都需要从磁盘读取数据并且在计算完成后需要将中间结果写入到磁盘中IO 开销较大。3延迟高一次计算可能需要分解成一系列按顺序执行的 MapReduce 任务任务之间的衔接由于涉及到 IO 开销会产生较高延迟。而且在前一个任务执行完成之前其他任务无法开始因此难以胜任复杂、多阶段的计算任务。由于 MapReduce 是基于磁盘的分布式计算框架因此性能方面要逊色于基于内存的分布式计算框比如 Spark 和 Flink 等所以随着 Spark 和 Flink 的发展MapReduce 的市场空间逐渐被挤压地位也逐渐被边缘化但是不可否认的是MapReduce 的一些优秀的设计思想还是在其他框架中得到了很好的继承。数据仓库 HiveHive 简介Hive 是一个构建于 Hadoop 顶层的数据仓库工具某种程度上可以看作是用户编程接口本身不存储和处理数据依赖分布式文件系统 HDFS 存储数据依赖分布式并行计算模型 MapReduce 处理数据定义了简单的类 SQL 查询语言——HiveQL用户可以通过编写的 HiveQL 语句运行 MapReduce 任务是一个可以提供有效、合理、直观组织和使用数据的模型。Hive 与 Hadoop 生态系统中其他组件的关系如下图所示1Hive 依赖于 HDFS 存储数据。HDFS 作为高可靠性的底层存储用来存储海量数据。2Hive 依赖于 MapReduce 处理数据。MapReduce 对这些海量数据进行处理实现高性能计算用 HiveQL 语句编写的处理逻辑最终均要转化为 MapReduce 任务来运行。3Pig可以作为Hive的替代工具。Pig 是一种数据流语言和运行环境适合用于 Hadoop 和 MapReduce 平台上查询半结构化数据集。常用于 ETL 过程的一部分即将外部数据装载到 Hadoop 集群中然后转换为用户期待的数据格式。4HBase 提供数据的实时访问。HBase 一个面向列的、分布式的、可伸缩的数据库它可以提供数据的实时访问功能而 Hive 只能处理静态数据主要是 BI 报表数据所以 HBase 与 Hive 的功能是互补的它实现了 Hive 不能提供功能。Hive 系统架构Hive 系统架构如下图所示Hive 系统架构用户接口模块包括 CLI、HWI、JDBC、ODBC、Thrift Server 等。CLI 是 Hive 自带的一个命令行界面HWI 是 Hive 的一个简单网页界面JDBC、ODBC 以及 Thrift Server 可以向用户提供进行编程访问的接口。驱动模块包括编译器、优化器、执行器等。所有命令和查询都会进入到驱动模块通过该模块对输入进行解析编译对需求的计算进行优化然后按照指定的步骤进行执行。元数据存储模块Metastore是一个独立的关系型数据库。通常是与 MySQL 数据库连接后创建的一个 MySQL 实例也可以是 Hive 自带的 derby 数据库实例。元数据存储模块中主要保存表模式和其他系统元数据如表的名称、表的列及其属性、表的分区及其属性、表的属性、表中数据所在位置信息等。数据仓库 ImpalaImpala 是由 Cloudera 公司开发的新型查询系统它提供 SQL 语义能查询存储在 Hadoop 的 HDFS 和 HBase 上的 PB 级大数据。Impala 最开始是参照 Dremel 系统进行设计的Impala 的目的不在于替换现有的 MapReduce 工具而是提供一个统一的平台用于实时查询。Impala 与其他组件关系如下图所示与 Hive 类似Impala 也可以直接与 HDFS 和 HBase 进行交互。Hive 底层执行使用的是 MapReduce所以主要用于处理长时间运行的批处理任务例如批量提取、转化、加载类型的任务。Impala 通过与商用并行关系数据库中类似的分布式查询引擎可以直接从 HDFS 或者 HBase 中用 SQL 语句查询数据从而大大降低了延迟主要用于实时查询。Impala 和 Hive 采用相同的 SQL 语法、ODBC 驱动程序和用户接口。基于内存的分布式计算框架 SparkSpark 简介Spark 最初由美国加州大学伯克利分校UC Berkeley的AMP实验室于 2009 年开发是基于内存计算的大数据并行计算框架可用于构建大型的、低延迟的数据分析应用程序。2013 年 Spark 加入 Apache 孵化器项目后发展迅猛如今已成为 Apache 软件基金会最重要的三大分布式计算系统开源项目之一Hadoop、Spark、Storm。Spark 在 2014 年打破了 Hadoop 保持的基准排序纪录Spark / 206 个节点 / 23分钟 / 100TB 数据Hadoop / 2000 个节点 / 72 分钟 / 100TB 数据Spark 用十分之一的计算资源获得了比 Hadoop 快 3 倍的速度。Spark 具有如下几个主要特点1运行速度快使用 DAG 执行引擎以支持循环数据流与内存计算。2容易使用支持使用 Scala、Java、Python 和 R 语言进行编程可以通过Spark Shell 进行交互式编程。3通用性Spark 提供了完整而强大的技术栈包括 SQL 查询、流式计算、机器学习和图算法组件。4运行模式多样可运行于独立的集群模式中可运行于 Hadoop 中也可运行于 Amazon EC2 等云环境中并且可以访问 HDFS、Cassandra、HBase、Hive 等多种数据源。Spark 相对于 MapReduce 的优点Spark 在借鉴 Hadoop MapReduce 优点的同时很好地解决了 MapReduce 所面临的问题。相比于 Hadoop MapReduceSpark 主要具有如下优点1Spark 的计算模式也属于 MapReduce但不局限于 Map 和 Reduce 操作还提供了多种数据集操作类型编程模型比 Hadoop MapReduce 更灵活.2Spark 提供了内存计算可将中间结果放到内存中对于迭代运算效率更高。3Spark 基于 DAG 的任务调度执行机制要优于 Hadoop MapReduce 的迭代执行机制。Hadoop MapReduce 与 Spark 的执行流程对比如下图所示Hadoop MapReduce 执行流程Spark 执行流程使用 Hadoop 进行迭代计算非常耗资源Spark 将数据载入内存后之后的迭代计算都可以直接使用内存中的中间结果作运算避免了从磁盘中频繁读取数据。 Hadoop 与 Spark 执行逻辑回归的时间对比如下图所示Spark 与 Hadoop 的关系Spark 正以其结构一体化、功能多元化的优势逐渐成为当今大数据领域最热门的大数据计算平台。目前越来越多的企业放弃 MapReduce转而使用 Spark 开发企业应用。但是需要指出的是Spark 作为计算框架只能解决数据计算问题无法解决数据存储问题Spark 只是取代了 Hadoop 生态系统中的计算框架 MapReduce而 Hadoop 中的其他组件依然在企业大数据系统中发挥着重要的作用。比如企业在采用 Spark 解决数据计算问题的同时依然需要依赖 Hadoop 分布式文件系统 HDFS 和分布式数据库 HBase来实现不同类型数据的存储和管理并借助于 YARN 实现集群资源的管理和调度。因此在许多企业实际应用中Hadoop 和 Spark 的统一部署是一种比较现实合理的选择。由于MapReduce、Storm 和 Spark等都可以运行在资源管理框架 YARN 之上因此可以在 YARN 之上统一部署各个计算框架。如下图所示Hadoop 和 Spark 的统一部署这些不同的计算框架统一运行在 YARN 中可以带来如下好处1计算资源按需伸缩2不用负载应用混搭集群利用率高3共享底层存储避免数据跨集群迁移。Spark 生态系统在实际应用中大数据处理主要包括以下三个类型1复杂的批量数据处理通常时间跨度在数十分钟到数小时之间2基于历史数据的交互式查询通常时间跨度在数十秒到数分钟之间3基于实时数据流的数据处理通常时间跨度在数百毫秒到数秒之间当同时存在以上三种场景时就需要同时部署三种不同的软件比如: MapReduce / Impala / Storm 。这样做难免会带来一些问题不同场景之间输入输出数据无法做到无缝共享通常需要进行数据格式的转换不同的软件需要不同的开发和维护团队带来了较高的使用成本比较难以对同一个集群中的各个系统进行统一的资源协调和分配。Spark 的设计遵循“一个软件栈满足不同应用场景”的理念逐渐形成了一套完整的生态系统既能够提供内存计算框架也可以支持 SQL 即席查询、实时流式计算、机器学习和图计算等。Spark 可以部署在资源管理器 YARN 之上提供一站式的大数据解决方案因此Spark 所提供的生态系统足以应对上述三种场景即同时支持批处理、交互式查询和流数据处理。Spark 的生态系统主要包含了 Spark Core、Spark SQL、Spark Streaming、Structured Streaming、MLlib 和 GraphX 等组件。如下图所示Spark 生态系统Spark 生态系统组件的应用场景应用场景 时间跨度 其他框架 Spark 生态系统中的组件复杂的批量数据处理 小时级 MapReduce、Hive Spark基于历史数据的交互式查询 分钟级、秒级 Impala、Dremel、Drill Spark SQL基于实时数据流的数据处理 毫秒、秒级 Storm、S4 Spark Streaming、Structured Streaming基于历史数据的数据挖掘 - Mahout MLlib图结构数据的处理 - Pregel、Hama GraphXSpark 体系架构Spark 运行架构包括集群资源管理器Cluster Manager、运行作业任务的工作节点Worker Node、每个应用的任务控制节点Driver和每个工作节点上负责具体任务的执行进程Executor资源管理器可以自带或 Mesos 或 YARN。Spark 运行架构Spark 的数据抽象 RDDSpark Core 是建立在统一的抽象 RDD 之上使得 Spark 的各个组件可以无缝进行集成在同一个应用程序中完成大数据计算任务。一个 RDD 就是一个分布式对象集合本质上是一个只读的分区记录集合每个 RDD 可以分成多个分区每个分区就是一个数据集片段并且一个 RDD 的不同分区可以被保存到集群中不同的节点上从而可以在集群中的不同节点上进行并行计算。RDD 提供了一组丰富的操作以支持常见的数据运算分为“动作”Action和“转换”Transformation两种类型RDD 提供的转换接口都非常简单都是类似 map、filter、groupBy、join 等粗粒度的数据转换操作而不是针对某个数据项的细粒度修改不适合网页爬虫。表面上 RDD 的功能很受限、不够强大实际上 RDD 已经被实践证明可以高效地表达许多框架的编程模型比如 MapReduce、SQL、PregelSpark 提供了 RDD 的 API程序员可以通过调用 API 实现对 RDD 的各种操作。RDD 典型的执行过程如下1RDD 读入外部数据源进行创建2RDD 经过一系列的转换Transformation操作每一次都会产生不同的RDD供给下一个转换操作使用3最后一个 RDD 经过“动作”操作进行转换并输出到外部数据源。Spark 的部署方式Spark 支持五种不同类型的部署方式包括Local、Standalone类似于 MapReduce1.0slot 为资源分配单位、Spark on Mesos和 Spark 有血缘关系更好支持 Mesos、Spark on YARN、Spark on Kubernetes。其中 Spark on YARN 的架构如下图所示Spark on YARN 架构Spark SQL1、从 Shark 说起Shark 即 Hive on Spark为了实现与 Hive 兼容Shark 在 HiveQL 方面重用了 Hive 中 HiveQL 的解析、逻辑执行计划翻译、执行计划优化等逻辑可以近似认为仅将物理执行计划从 MapReduce 作业替换成了 Spark 作业通过 Hive 的 HiveQL 解析把 HiveQL 翻译成 Spark 上的 RDD 操作。Shark 直接继承了 Hive 的各个组件Shark 的出现使得 SQL-on-Hadoop 的性能比 Hive 有了 10-100 倍的提高。Shark 的设计导致了两个问题1一是执行计划优化完全依赖于 Hive不方便添加新的优化策略。2二是因为 Spark 是线程级并行而 MapReduce 是进程级并行因此Spark 在兼容 Hive 的实现上存在线程安全问题导致 Shark 不得不使用另外一套独立维护的打了补丁的 Hive 源码分支。2014 年 6 月 1 日 Shark 项目和 Spark SQL 项目的主持人 Reynold Xin 宣布停止对 Shark 的开发团队将所有资源放在 Spark SQL 项目上至此Shark 的发展画上了句号但也因此发展出两个分支Spark SQL 和 Hive on Spark。Spark SQL 作为 Spark 生态的一员继续发展而不再受限于 Hive只是兼容 HiveHive on Spark 是一个 Hive 的发展计划该计划将 Spark 作为 Hive 的底层引擎之一也就是说Hive 将不再受限于一个引擎可以采用 Map-Reduce、Tez、Spark 等引擎。2、Spark SQL 架构Spark SQL 在 Hive 兼容层面仅依赖 HiveQL 解析、Hive 元数据也就是说从 HQL 被解析成抽象语法树AST起就全部由 Spark SQL 接管了。Spark SQL 执行计划生成和优化都由 Catalyst函数式关系查询优化框架负责。Spark SQL 架构Spark SQL 增加了 DataFrame即带有 Schema 信息的 RDD使用户可以在 Spark SQL 中执行 SQL 语句数据既可以来自 RDD也可以是 Hive、HDFS、Cassandra 等外部数据源还可以是 JSON 格式的数据。Spark SQL 目前支持 Scala、Java、Python 三种语言支持 SQL-92 规范Spark SQL 支持的数据格式和编程语言3、为什么推出 Spark SQL1关系数据库已经很流行2关系数据库在大数据时代已经不能满足要求首先用户需要从不同数据源执行各种操作包括结构化、半结构化和非结构化数据其次用户需要执行高级分析比如机器学习和图像处理在实际大数据应用中经常需要融合关系查询和复杂分析算法比如机器学习或图像处理但是缺少这样的系统。Spark SQL 填补了这个鸿沟首先可以提供 DataFrame API可以对内部和外部各种数据源执行各种关系型操作其次可以支持大数据中的大量数据源和数据分析算法 Spark SQL 可以融合传统关系数据库的结构化数据管理能力和机器学习算法的数据处理能力。Spark StreamingSpark Streaming 是构建在 Spark Core 上的实时计算框架它扩展了 Spark 处理大规模流式数据的能力。Spark Streaming 可结合批处理和交互查询适合一些需要对历史数据和实时数据进行结合分析的应用场景。Spark Streaming 是 Spark 的核心组件之一为 Spark 提供了可拓展、高吞吐、容错的流计算能力。Spark Streaming 可整合多种输入数据源如 Kafka、Flume、HDFS甚至是普通的 TCP 套接字如下图所示。Spark Streaming 支持的输入、输出数据源经处理后的数据可存储至文件系统、数据库或显示在仪表盘里。Spark Streaming 实际上是以一系列微小批处理来模拟流计算。Spark Streaming 的基本原理是将实时输入数据流以时间片秒级为单位进行拆分然后经 Spark 引擎以类似批处理的方式处理每个时间片数据。执行流程如下图所示Spark Streaming 执行流程下图展示了Spark Streaming 的工作机制。Spark Streaming 工作机制在 Spark Streaming 中会有一个组件 Receiver作为一个长期运行的任务task跑在一个 Executor 上每个 Receiver 都会负责一个 DStream 输入流比如从文件中读取数据的文件流、套接字流或者从 Kafka 中读取的一个输入流等等。Receiver 组件接收到数据源发来的数据后会提交给 Spark Streaming 程序进行处理。处理后的结果可以交给可视化组件进行可视化展示或者也可以写入到 HDFS、HBase 中。Structured Streaming1、Structured Streaming 简介Structured Streaming 是一种基于 Spark SQL 引擎构建的、可扩展且容错的流处理引擎。通过一致的 APIStructured Streaming 使得使用者可以像写批处理程序一样编写流处理程序简化了使用者的使用难度。提供端到端的完全一致性是 Structured Streaming 设计背后的关键目标之一为了实现这一点Spark 设计了输入源、执行引擎和接收器以便对处理的进度进行更可靠地跟踪使之可以通过重启或重新处理来处理任何类型的故障。如果所使用的源具有偏移量来跟踪流的读取位置那么引擎可以使用检查点和预写日志来记录每个触发时期正在处理的数据的偏移范围此外如果使用的接收器是“幂等”的那么通过使用重放、对“幂等”接收数据进行覆盖等操作Structured Streaming 可以确保在任何故障下达到端到端的完全一致性。Spark 一直在不停更新中从 Spark 2.3.0 版本开始引入了持续流式处理模型可以将原先流处理的延迟降低到毫秒级别。2、Structured Streaming 的关键思想Structured Streaming 的关键思想是将实时数据流视为一张正在不断添加数据的表可以把流计算等同于在一个静态表上的批处理查询Spark 会在不断添加数据的无界输入表上运行计算并进行增量查询。如下图所示无界表在无界表上对输入的查询将生成结果表系统每隔一定的周期会触发对无界表的计算并更新结果表。如下图所示Structured Streaming 编程模型3、Structured Streaming 的两种处理模型1微批处理Structured Streaming 默认使用微批处理执行模型这意味着 Spark 流计算引擎会定期检查流数据源并对自上一批次结束后到达的新数据执行批量查询数据到达和得到处理并输出结果之间的延时超过 100 毫秒。Structured Streaming 的微批处理模型2持续处理Spark 从 2.3.0 版本开始引入了持续处理的试验性功能可以实现流计算的毫秒级延迟在持续处理模式下Spark 不再根据触发器来周期性启动任务而是启动一系列的连续读取、处理和写入结果的长时间运行的任务。4、Structured Streaming 和 Spark SQL、Spark Streaming 关系Structured Streaming 处理的数据跟 Spark Streaming 一样也是源源不断的数据流区别在于Spark Streaming 采用的数据抽象是 DStream本质上就是一系列 RDD而 Structured Streaming 采用的数据抽象是 DataFrame。Structured Streaming 可以使用 Spark SQL 的 DataFrame/Dataset 来处理数据流。虽然 Spark SQL 也是采用 DataFrame 作为数据抽象但是SparkSQL 只能处理静态的数据而 Structured Streaming 可以处理结构化的数据流。这样Structured Streaming 就将 Spark SQL 和 Spark Streaming 二者的特性结合了起来。Structured Streaming 可以对 DataFrame/Dataset 应用前面提到的各种操作包括 select、where、groupBy、map、filter、flatMap 等。Spark Streaming 只能实现秒级的实时响应而 Structured Streaming 由于采用了全新的设计方式采用微批处理模型时可以实现 100 毫秒级别的实时响应采用持续处理模型时可以支持毫秒级的实时响应。Spark MLlib传统的机器学习算法由于技术和单机存储的限制只能在少量数据上使用依赖于数据抽样大数据技术的出现可以支持在全量数据上进行机器学习机器学习算法涉及大量迭代计算基于磁盘的 MapReduce 不适合进行大量迭代计算而基于内存的 Spark 比较适合进行大量迭代计算。Spark 提供了一个基于海量数据的机器学习库它提供了常用机器学习算法的分布式实现。开发者只需要有 Spark 基础并且了解机器学习算法的原理以及方法相关参数的含义就可以轻松的通过调用相应的 API 来实现基于海量数据的机器学习过程。pyspark 的即席查询也是一个关键算法工程师可以边写代码边运行边看结果。MLlib 是 Spark 的机器学习Machine Learning库旨在简化机器学习的工程实践工作MLlib 由一些通用的学习算法和工具组成包括分类、回归、聚类、协同过滤、降维等同时还包括底层的优化原语和高层的流水线PipelineAPI具体如下算法工具常用的学习算法如分类、回归、聚类和协同过滤特征化工具特征提取、转化、降维和选择工具流水线Pipeline用于构建、评估和调整机器学习工作流的工具持久性保存和加载算法、模型和管道实用工具线性代数、统计、数据处理等工具。TensorFlowOnSparkTensorFlow 是一个开源的、基于 Python 的机器学习框架它是由谷歌公司开发的并在图形分类、音频处理、推荐系统和自然语言处理等场景下有着丰富的应用是目前最热门的机器学习框架。TensorFlow 是一个采用数据流图Data Flow Graph、用于数值计算的开源软件库。数据流图中的节点Nodes表示数学操作图中的线则表示节点间的相互联系的多维数据组即张量 (Tensor)。在计算过程中张量从图的一端流动到另一端这也是这个工具取名为“TensorFlow”的原因。一旦输入端的所有张量准备好节点将被分配到各种计算设备完成异步并行地执行运算。利用 TensorFlow 我们可以在多种平台上展开数据分析与计算如 CPU或 GPU、台式机、服务器、甚至移动设备等等。尽管 TensorFlow 也开放了自己的分布式运行框架但在目前公司的技术架构和使用环境上不是那么友好如何将 TensorFlow 加入到现有的环境中Spark/YARN并为用户提供更加方便易用的环境成为了目前所要解决的问题。TensorFlowOnSpark 项目是由 Yahoo 开源的一个软件包能将 TensorFlow 与 Spark 结合在一起使用为 Apache Hadoop 和 Apache Spark 集群带来可扩展的深度学习功能。使 Spark 能够利用 TensorFlow 拥有深度学习和 GPU 加速计算的能力。传统情况下处理数据需要跨集群深度学习集群和 Hadoop/Spark 集群Yahoo 为了解决跨集群传递数据的问题开发了 TensorFlowOnSpark 项目。TensorFlowOnSpark 目前被用于雅虎私有云中的 Hadoop 集群主要进行大规模分布式深度学习。TensorFlowOnSpark 在设计时充分考虑了 Spark 本身的特性和 TensorFlow 的运行机制大大保证了两者的兼容性使得可以通过较少的修改来运行已经存在的 TensorFlow 程序。在独立的 TFOnSpark 程序中能够与 SparkSQL、MLlib 和其他 Spark 库一起工作处理数据如下图所示TensorFlowOnSpark 与 Spark 的集成TensorFlowOnSpark 的体系架构较为简单如图所示Spark Driver 程序并不会参与 TensorFlow 内部相关的计算和处理。其设计思路像是将一个 TensorFlow 集群运行在了 Spark 上它会在每个 Spark Executor 中启动 TensorFlow 应用程序然后通过 gRPC 或 RDMA 方式进行数据传递与交互。TensorFlowOnSpark 体系架构TensorFlowOnSpark 的 Spark 应用程序包括4 个基本过程1预留组建 TensorFlow 集群并在每个 Executor 进程上预留监听端口启动“数据/控制”消息的监听程序2启动在每个 Executor 进程上启动 TensorFlow 应用程序3训练/推理在 TensorFlow 集群上完成模型的训练或推理4关闭关闭 Executor 进程上的 TensorFlow 应用程序释放相应的系统资源(消息队列)。流计算框架 StormStorm 简介Twitter Storm 是一个免费、开源的分布式实时计算系统Storm 对于实时计算的意义类似于 Hadoop 对于批处理的意义Storm 可以简单、高效、可靠地处理流数据并支持多种编程语言Storm 框架可以方便地与数据库系统进行整合从而开发出强大的实时计算系统。Twitter 是全球访问量最大的社交网站之一Twitter 开发 Storm 流处理框架也是为了应对其不断增长的流数据实时处理需求。如下图所示Twitter 的分层数据处理架构Storm 的特点Storm 可用于许多领域中如实时分析、在线机器学习、持续计算、远程 RPC、数据提取加载转换等Storm 具有以下主要特点1整合性Storm 可方便地与队列系统和数据库系统进行整合2简易的APIStorm 的 API 在使用上即简单又方便3可扩展性Storm 的并行特性使其可以运行在分布式集群中4容错性Storm 可自动进行故障节点的重启、任务的重新分配5可靠的消息处理Storm 保证每个消息都能完整处理6支持各种编程语言Storm 支持使用各种编程语言来定义任务7快速部署Storm 可以快速进行部署和使用8免费、开源Storm 是一款开源框架可以免费使用。Storm 的框架设计Storm 运行任务的方式与 Hadoop 类似Hadoop 运行的是 MapReduce 作业而 Storm 运行的是“Topology”但两者的任务大不相同主要的不同是MapReduce 作业最终会完成计算并结束运行而 Topology 将持续处理消息直到人为终止。Storm 和 Hadoop 架构组件功能对应关系Hadoop Storm应用名称 Job Topology系统角色 JobTracker NimbusTaskTracker Supervisor组件接口 Map/Reduce Spout/BoltStorm 集群采用“Master—Worker”的节点方式Master 节点运行名为“Nimbus”的后台程序类似 Hadoop 中的“JobTracker”负责在集群范围内分发代码、为 Worker 分配任务和监测故障。Worker 节点运行名为“Supervisor”的后台程序负责监听分配给它所在机器的工作即根据 Nimbus 分配的任务来决定启动或停止 Worker 进程一个 Worker 节点上同时运行若干个 Worker 进程。Storm 使用 Zookeeper 来作为分布式协调组件负责 Nimbus 和多个 Supervisor 之间的所有协调工作。借助于 Zookeeper若 Nimbus 进程或 Supervisor 进程意外终止重启时也能读取、恢复之前的状态并继续工作使得 Storm 极其稳定。Storm 集群架构示意图基于这样的架构设计Storm 的工作流程如下图所示Storm 工作流程示意图所有 Topology 任务的提交必须在 Storm 客户端节点上进行提交后由 Nimbus 节点分配给其他 Supervisor 节点进行处理Nimbus 节点首先将提交的 Topology 进行分片分成一个个 Task分配给相应的 Supervisor并将 Task 和 Supervisor 相关的信息提交到 Zookeeper 集群上。Supervisor 会去 Zookeeper 集群上认领自己的 Task通知自己的 Worker 进程进行 Task 的处理。在提交了一个 Topology 之后Storm 就会创建 Spout/Bolt 实例并进行序列化。之后将序列化的组件发送给所有的任务所在的机器(即 Supervisor 节点)在每一个任务上反序列化组件。Spark Streaming 与 Storm 的对比1Spark Streaming 和 Storm 最大的区别在于Spark Streaming 无法实现毫秒级的流计算而 Storm 可以实现毫秒级响应。2Spark Streaming 构建在 Spark 上一方面是因为 Spark 的低延迟执行引擎100ms可以用于实时计算另一方面相比于 StormRDD 数据集更容易做高效的容错处理。3Spark Streaming 采用的小批量处理的方式使得它可以同时兼容批量和实时数据处理的逻辑和算法因此方便了一些需要历史数据和实时数据联合分析的特定应用场合。从编程的灵活性来讲Storm 是比较理想的选择它使用 ApacheThrift可以用任何编程语言来编写拓扑结构Topology当需要在一个集群中把流计算和图计算、机器学习、SQL 查询分析等进行结合时可以选择 Spark Streaming因为在 Spark 上可以统一部署 Spark SQLSpark Streaming、MLlibGraphX 等组件提供便捷的一体化编程模型当应用场景需要毫秒级响应时可以选择 Storm因为 Spark Streaming 无法实现毫秒级的流计算。流计算框架 FlinkFlink 简介Flink 是 Apache 软件基金会的顶级项目之一是一个针对流数据和批数据的分布式计算框架设计思想主要来源于 Hadoop、MPP 数据库、流计算系统等。Flink 主要是由 Java 代码实现的目前主要还是依靠开源社区的贡献而发展。Flink 所要处理的主要场景是流数据批数据只是流数据的一个特例而已也就是说Flink 会把所有任务当成流来处理。Flink 可以支持本地的快速迭代以及一些环形的迭代任务。Flink 以层级式系统形式组建其软件栈如下图所示Flink 架构图不同层的栈建立在其下层基础上。具体而言Flink 的典型特性如下1提供了面向流处理的 DataStream API 和面向批处理的 DataSet API。DataSet API 支持 Java、Scala 和 PythonDataStream API 支持 Java 和 Scala2提供了多种候选部署方案比如本地模式Local、集群模式Cluster和云模式Cloud。对于集群模式而言可以采用独立模式Standalone或者YARN3提供了一些类库包括 Table处理逻辑表查询、FlinkML机器学习、Gelly图像处理和 CEP复杂事件处理4提供了较好的 Hadoop 兼容性不仅可以支持 YARN还可以支持 HDFS、HBase 等数据源。Flink 和 Spark 一样都是基于内存的计算框架因此都可以获得较好的实时计算性能。当全部运行在 Hadoop YARN 之上时Flink 的性能甚至还要略好于 Spark因为Flink 支持增量迭代具有对迭代进行自动优化的功能。Flink 和Spark Streaming 都支持流计算二者的区别在于Flink 是一行一行地处理数据而 Spark Streaming 是基于 RDD 的小批量处理所以Spark Streaming 在流式处理方面不可避免地会增加一些延时实时性没有 Flink 好。Flink 的流计算性能和 Storm 差不多可以支持毫秒级的响应而 Spark Streaming 则只能支持秒级响应。总体而言Flink 和 Spark 都是非常优秀的基于内存的分布式计算框架但是Spark 的市场影响力和社区活跃度明显超过 Flink这在一定程度上限制了 Flink 的发展空间。Flink 是理想的流计算框架流处理架构需要具备低延迟、高吞吐和高性能的特性而目前从市场上已有的产品来看只有 Flink 可以满足要求。Storm 虽然可以做到低延迟但是无法实现高吞吐也不能在故障发生时准确地处理计算状态。Spark Streaming 通过采用微批处理方法实现了高吞吐和容错性但是牺牲了低延迟和实时处理能力。Structured Streaming 采用持续处理模型时可以支持毫秒级的实时响应但是这是以牺牲一致性为代价的持续处理模型只能做到“至少一次”的一致性而无法保证端到端的完全一致性。Flink 实现了 Google Dataflow 流计算模型是一种兼具高吞吐、低延迟和高性能的实时流计算框架并且同时支持批处理和流处理。此外Flink 支持高度容错的状态管理防止状态在计算过程中因为系统异常而出现丢失。因此Flink 就成为了能够满足流处理架构要求的理想的流计算框架。Flink 体系架构如下图所示Flink 系统主要由两个组件组成分别为 JobManager 和 TaskManagerFlink 架构也遵循 Master-Slave 架构设计原则JobManager 为 Master 节点TaskManager 为 Slave 节点。Flink 体系架构大数据编程框架 Beam在大数据处理领域开发者经常要用到很多不同的技术、框架、API、开发语言和 SDK。根据不同的企业业务系统开发需求开发者很可能会用 MapReduce 进行批处理用 Spark SQL 进行交互式查询用 Flink 实现实时流处理还有可能用到基于云端的机器学习框架。大量的开源大数据产品比如 MapReduce、Spark、Flink、Storm、Apex 等为大数据开发者提供了丰富的工具的同时也增加了开发者选择合适工具的难度尤其对于新入行的开发者来说更是如此。新的分布式处理框架可能带来更高的性能、更强大的功能和更低的延迟但是用户切换到新的分布式处理框架的代价也非常大——需要学习一个新的大数据处理框架并重写所有的业务逻辑。解决这个问题的思路包括两个部分首先需要一个编程范式能够统一、规范分布式数据处理的需求例如统一批处理和流处理的需求其次生成的分布式数据处理任务应该能够在各个分布式执行引擎如 Spark、Flink 等上执行用户可以自由切换分布式数据处理任务的执行引擎与执行环境。ApacheBeam 的出现就是为了解决这个问题。Beam 是由谷歌贡献的 Apache 顶级项目它的目标是为开发者提供一个易于使用、却又很强大的数据并行处理模型能够支持流处理和批处理并兼容多个运行平台。Beam 是一个开源的统一的编程模型开发者可以使用 Beam SDK 来创建数据处理管道然后这些程序可以在任何支持的执行引擎上运行比如运行在 Apex、Spark、Flink、Cloud Dataflow 上。Beam SDK 定义了开发分布式数据处理任务业务逻辑的 API 接口即提供一个统一的编程接口给到上层应用的开发者开发者不需要了解底层的具体的大数据平台的开发接口是什么直接通过 BeamSDK 的接口就可以开发数据处理的加工流程不管输入是用于批处理的有限数据集还是用于流处理的无限数据集。对于有限或无限的输入数据Beam SDK 都使用相同的类来表现并且使用相同的转换操作进行处理。如下图所示Beam 使用一套高层抽象的 API 屏蔽多种计算引擎的区别终端用户用 Beam 来实现自己所需的流计算功能使用的终端语言可能是 Python、Java 等Beam 为每种语言提供了一个对应的 SDK用户可以使用相应的 SDK 创建数据处理管道用户写出的程序可以被运行在各个 Runner 上每个 Runner 都实现了从 Beam 管道到平台功能的映射。目前主流的大数据处理框架 Flink、Spark、Apex 以及谷歌的 Cloud DataFlow 等都有了支持 Beam 的 Runner。通过这种方式Beam 使用一套高层抽象的 API 屏蔽了多种计算引擎的区别开发者只需要编写一套代码就可以运行在不同的计算引擎之上比如 Apex、Spark、Flink、Cloud Dataflow 等。查询分析系统 DremelDremel 是一种可扩展的、交互式的实时查询系统用于只读嵌套数据的分析。通过结合多级树状执行过程和列式数据结构它能做到几秒内完成对万亿张表的聚合查询。系统可以扩展到成千上万的 CPU 上满足谷歌公司上万用户操作 PB 级的数据可以在 2 到 3 秒内完成PB级别数据的查询。Dremel 具有以下几个主要的特点1Dremel 是一个大规模、稳定的系统。2Dremel 是 MapReduce 交互式查询能力不足的补充。3Dremel 的数据模型是嵌套的。4Dremel 中的数据是用列式存储的。5Dremel 结合了 Web 搜索和并行 DBMSDatabase Management System的技术。编程要求根据相关知识完成右侧选择题有任何问题都可以随时关注私信

相关新闻

2026/8/27 21:09:47

视频解码器接入MIPI-CSI2接口:从设备树配置到调试避坑指南

最近在做一块视频解码相关的板卡,核心任务是把一路视频源解码后,通过MIPI-CSI2接口送给新一代SoC处理。这块工作在前期沟通时看着不难,实际落地时从协议理解到信号调试,再到SoC端设备树配置,每一步都有不少讲究。这篇文…

2026/8/27 21:04:47

Jetson TX2/Xavier加固电脑Linux系统部署:从刷机到性能调优全指南

前阵子在一个户外项目中折腾 Jetson TX2 和 Xavier 平台的加固电脑,天天泡在刷机、调系统和压测里。这个领域其实很有意思:外人看着就是一台“防尘防水防摔的电脑”,但对搞嵌入式 Linux 的人来说,它是一整套和普通 PC 完全不同的技…

2026/8/27 21:04:47

单片机毕业设计-基于 STM32 的指纹刷卡密码一体化门禁装置开发 基于 STM32 单片机的多模态身份识别门禁系统研究(012505)

博主介绍:✌️码农一枚 ,专注于大学生项目实战开发、讲解和毕业🚢文撰写修改等。全栈领域优质创作者,博客之星、掘金/华为云/阿里云/InfoQ等平台优质作者、专注于嵌入式单片机,Java、小程序技术领域和毕业项目实战 ✌️…

2026/8/27 21:44:50

Rescene:免Key AI Agent聚合器的本地部署与使用指南

这次我们来看一个对本地部署和 AI Agent 折腾党很有用的开源项目:Rescene。它的定位很简单:一个免费的 AI Agent 聚合器,而且不需要用户提供 API Key。也就是说,你不需要先去某个大模型平台申请密钥、配置支付方式,再回…

2026/8/27 21:44:50

从OpenAI到Meta:JAX与强化学习趋势下的AI Agent实战

这两天技术圈有一则人事变动消息吸引了不少人注意:知名 AI 研究员 Luke Metz 被曝将离开 OpenAI,加入 Meta 的超级智能实验室。单看名字,很多人可能不熟悉,但提到 JAX、强化学习、大规模模型训练这些方向,Luke Metz 的…

2026/8/27 21:44:50

TrueTouch触控平台实战解析:从电容检测原理到系统调校

做触摸方案这些年,“TrueTouch”这个名字从我入行起就没离开过视野。赛普拉斯半导体(现在算英飞凌麾下)推出的这一系列电容式触摸屏控制器,曾经是手机上高端触控方案的代名词。到现在,很多平板、车载屏幕、智能家电上还…

2026/8/27 21:44:50

果园机器人视觉鲁棒性设计:从图像识别到工业落地

1. 这道赛题不是在考“识别苹果”,而是在考“如何让机器人在真实果园里不犯错” 2023年亚太数学建模竞赛A题,标题写着“水果采摘机器人的图像识别技术”,但如果你真把它当成一道普通的图像分类练习题来刷——比如用ResNet跑通一个苹果/香蕉/橙…

2026/8/27 21:39:49

AI大厨系统架构与部署实践:从视觉识别到批量出餐全拆解

3分钟出餐、30秒一杯咖啡。这组数字最近在餐饮和自动化圈里反复出现,不少朋友问我:这到底是营销话术,还是真的能落地的系统? 先给结论:这不是一台机器能单独完成的,而是一整套“AI 大厨”方案在起作用。它…

2026/8/26 9:13:28

[光学原理与应用-521]:对光的错误理解与纠偏

首先光是一种能量的载体和形态,宏观上观察到的光是由无数个微观的光量子组成的,每个光子在产生的瞬间,其在真空的空间中以确定不变的速度沿着一个初始的方向一直向前,在微观层面,每个光量子的运动轨迹是以波函数所展现…

2026/8/27 10:58:22

SIP通话转接原理与REFER方法实战解析

1. 通话转接不是“挂断再拨号”,而是SIP会话的动态重定向你有没有遇到过这样的场景:客服坐席A正在和客户通电话,突然需要把这通对话无缝转给专家坐席B,客户完全感知不到中间的断连——既没听到忙音,也没被要求重新拨号…

2026/8/27 7:46:21

Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

1. 为什么选择Kolla-ansible来部署单节点OpenStack?如果你正在寻找一种能把OpenStack从“概念”快速变成“可用的实验环境”的方法,那么Kolla-ansible几乎是当前最主流、最省心的选择。我见过太多人卡在手动编译依赖、配置服务、处理版本冲突的泥潭里&am…

2026/8/27 0:01:16

Go语言构建企业级AI服务网关:统一管理英伟达等AI接口调用

1. 项目概述:从零构建一个企业级的AI服务网关 最近在帮一个做内容审核的团队做技术架构升级,他们原来的业务里,每天有几十万张图片和短视频需要过审,最初是接了几个开源的AI模型自己部署,但效果和性能一直不太稳定。后…

2026/8/27 0:01:16

LeetCode Hot100(51-60)算法精解与面试技巧

1. 题目背景与核心价值"hot100(51-60)"这个标题看起来像是某个编程题库或算法练习集中的一组题目编号。在技术社区中,类似命名通常指向LeetCode、牛客网等平台的热门题目集合。作为刷过300题的算法老手,我理解这类题目的核心价值在于&#xff…

2026/8/27 0:01:16

CRC校验实战:从模2除法到HJ212协议排错

1. 为什么一个“校验码”能扛住工业现场90%的数据 corruption? 你有没有遇到过这样的场景:嵌入式设备通过RS-485上传温湿度数据,上位机偶尔收到一帧乱码——温度显示成-273℃,湿度跳到999%,但串口波形看起来完全正常&a…

2026/8/26 19:34:06

实测才敢推 AI论文网站 2026最新测评与推荐

2026年真正好用的AI论文网站,核心看生成的论文质量、低AI味、格式正确、学术适配四大指标。综合实测,千笔AI、ThouPen、豆包、DeepSeek、Grammarly 是当前最值得推荐的梯队,覆盖从免费到付费、从中文到英文、从文科到理工的全场景需求。一、综…

2026/8/26 19:17:08

2026必备!AI论文网站测评:最新推荐与深度对比

2026年真正好用的AI论文网站,核心看生成的论文质量、低AI味、格式正确、学术适配四大指标。综合实测,千笔AI、ThouPen、豆包、DeepSeek、Grammarly 是当前最值得推荐的梯队,覆盖从免费到付费、从中文到英文、从文科到理工的全场景需求。 一、…

2026/8/26 19:34:05

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

一天写完毕业论文在2026年已不再是天方夜谭。2026年最炸裂、实测能大幅提速的AI论文写作工具,覆盖选题构思、文献整理、内容生成、格式排版等核心场景,真正帮你高效搞定论文难题。 一、全流程王者:一站式搞定论文全链路(一天定稿首…