Apache Beam: 下一代的大数据处理标准

本文涉及的产品
实时计算 Flink 版,5000CU*H 3个月
简介:

Apache Beam(原名Google DataFlow)是Google在2016年2月份贡献给Apache基金会的Apache孵化项目,被认为是继MapReduce,GFS和BigQuery等之后,Google在大数据处理领域对开源社区的又一个非常大的贡献。Apache Beam的主要目标是统一批处理和流处理的编程范式,为无限,乱序,web-scale的数据集处理提供简单灵活,功能丰富以及表达能力十分强大的SDK。Apache Beam项目重点在于数据处理的编程范式和接口定义,并不涉及具体执行引擎的实现,Apache Beam希望基于Beam开发的数据处理程序可以执行在任意的分布式计算引擎上。本文主要介绍Apache Beam的编程范式-Beam Model,以及通过Beam SDK如何方便灵活的编写分布式数据处理业务逻辑,希望读者能够通过本文对Apache Beam有初步的了解,同时对于分布式数据处理系统如何处理乱序无限数据流的能力有初步的认识。

Apache Beam基本架构

随着分布式数据处理不断发展,新的分布式数据处理技术也不断被提出,业界涌现出了越来越多的分布式数据处理框架,从最早的Hadoop MapReduce,到Apache Spark,Apache Storm,以及更近的Apache Flink,Apache Apex等。新的分布式处理框架可能带来的更高的性能,更强大的功能,更低的延迟等,但用户切换到新的分布式处理框架的代价也非常大:需要学习一个新的数据处理框架,并重写所有的业务逻辑。解决这个问题的思路包括两个部分,首先,需要一个编程范式,能够统一,规范分布式数据处理的需求,例如,统一批处理和流处理的需求。其次,生成的分布式数据处理任务应该能够在各个分布式执行引擎上执行,用户可以自由切换分布式数据处理任务的执行引擎与执行环境。Apache Beam正是为了解决以上问题而提出的。

Apache Beam主要由Beam SDK和Beam Runner组成,Beam SDK定义了开发分布式数据处理任务业务逻辑的API接口,生成的的分布式数据处理任务Pipeline交给具体的Beam Runner执行引擎。Apache Beam目前支持的API接口是由Java语言实现的,Python版本的API正在开发之中。Apache Beam支持的底层执行引擎包括Apache Flink,Apache Spark以及Google Cloud Platform,此外Apache Storm,Apache Hadoop,Apache Gearpump等执行引擎的支持也在讨论或开发当中。其基本架构如下图所示:


图1 Apache Beam架构图

需要注意的是,虽然Apache Beam社区非常希望所有的Beam执行引擎都能够支持Beam SDK定义的功能全集,但是在实际实现中可能并不一定。例如,基于MapReduce的Runner显然很难实现和流处理相关的功能特性。目前Google DataFlow Cloud是对Beam SDK功能集支持最全面的执行引擎,在开源执行引擎中,支持最全面的则是Apache Flink。

Beam Model

Beam Model指的是Beam的编程范式,即Beam SDK背后的设计思想。在介绍Beam Model之前,先简要介绍一下Beam Model要处理的问题域与一些基本概念。

  1. 数据。分布式数据处理要处理的数据类型一般可以分为两类,有限的数据集和无限的数据流。有限的数据集,比如一个HDFS中的文件,一个HBase表等,特点是数据提前已经存在,一般也已经持久化,不会突然消失。而无限的数据流,比如kafka中流过来的系统日志流,或是从twitter API拿到的twitter流等等,这类数据的特点是,数据动态流入,无穷无尽,无法全部持久化。一般来说,批处理框架的设计目标是用来处理有限的数据集,流处理框架的设计目标是用来处理无限的数据流。有限的数据集可以看做是无限的数据流的一种特例,但是从数据处理逻辑的角度,这两者并无不同之处,例如,假设微博数据包含时间戳和转发量,用户希望按照统计每小时的转发量总和,此业务逻辑应该可以同时在有限数据集和无限数据流上执行,并不应该因为数据源的不同而对业务逻辑的实现产生任何影响。
  2. 时间。Process Time是指数据进入分布式处理框架的时间,而Event-Time则是指数据产生的时间。这两个时间通常是不同的,例如,对于一个处理微博数据的流计算任务,一条2016-06-01-12:00:00发表的微博经过网络传输等延迟可能在2016-06-01-12:01:30才进入到流处理系统中。批处理任务通常进行全量的数据计算,较少关注数据的时间属性,但是对于流处理任务来说,由于数据流是无情无尽的,无法进行全量的计算,通常是对某个窗口中得数据进行计算,对于大部分的流处理任务来说,按照时间进行窗口划分,可能是最常见的需求。
  3. 乱序。对于流处理框架处理的数据流来说,其数据的到达顺序可能并不严格按照Event-Time的时间顺序。如果基于Process Time定义时间窗口,数据到达的顺序就是数据的顺序,因此不存在乱序问题。但是对于基于Event Time定义的时间窗口来说,可能存在时间靠前的消息在时间靠后的消息后到达的情况,这在分布式的数据源中可能非常常见。对于这种情况,如何确定迟到数据,以及对于迟到数据如何处理通常是很棘手的问题。

Beam Model处理的目标数据是无限的时间乱序数据流,不考虑时间顺序或是有限的数据集可看做是无限乱序数据流的一个特例。Beam Model从下面四个维度归纳了用户在进行数据处理的时候需要考虑的问题:

  1. What。如何对数据进行计算?例如,Sum,Join或是机器学习中训练学习模型等。在Beam SDK中由Pipeline中的操作符指定。
  2. Where。数据在什么范围中计算?例如,基于Process-Time的时间窗口,基于Event-Time的时间窗口,滑动窗口等等。在BeamSDK中由Pipeline中的窗口指定。
  3. When。何时将计算结果输出?例如,在1小时的Event-Time时间窗口中,每隔1分钟,将当前窗口计算结果输出。在Beam SDK中由Pipeline中的Watermark和触发器指定。
  4. How。迟到数据如何处理?例如,将迟到数据计算增量结果输出,或是将迟到数据计算结果和窗口内数据计算结果合并成全量结果输出。在Beam SDK中由Accumulation指定。

Beam Model将”WWWH“四个维度抽象出来组成了Beam SDK,用户在基于Beam SDK构建数据处理业务逻辑时,在每一步只需要根据业务需求按照这四个维度调用具体的API即可生成分布式数据处理Pipeline,并提交到具体执行引擎上执行。“WWWH”四个维度的抽象仅仅关注业务逻辑本身,和分布式任务如何执行没有任何关系。

Beam SDK

不同于Apache Flink或是Apache Spark,Beam SDK使用同一套API表示数据源,输出目标以及操作符等。下面介绍4个基于Beam SDK的数据处理任务,通过这四个数据处理任务,读者可以了解通过Beam Mode是如何统一灵活的描述批处理和流处理任务的,这4个任务用来处理手机游戏领域的统计需求,包括:

  1. 用户分数。批处理任务,基于有限数据集统计用户分数。
  2. 每小时团队分数。批处理任务,基于有限数据集统计每小时,每个团队的分数。
  3. 排行榜。流处理任务,2个统计项,每小时每个团队的分数以及用户实时的历史总得分数。
  4. 游戏状态。流处理任务,统计每小时每个团队的分数,以及更复杂的每小时统计信息,比如每小时每个用户在线时间等。

注:示例代码来自Beam的源码,具体地址参见:apache/incubator-beam。部分分析内容参考了Beam的官方文档,详情请参见引用链接。

下面基于Beam Model的“WWWH”四个维度,分析业务逻辑,并通过代码展示如何通过Beam SDK实现“WWWH”四个维度的业务逻辑。

用户分数

统计每个用户的历史总得分数是一个非常简单的任务,在这里我们简单的通过一个批处理任务实现,每次需要新的用户分数数据的时候,重新执行一次这个批处理任务即可。对于用户分数任务,“WWWH”四维度分析结果如下:

通过“WWWH”的分析,对于用户分数这个批处理任务,通过Beam Java SDK实现的代码如下所示:

 
  1. gameEvents  
  2. [... input ...]  
  3. [... parse ...]  
  4. .apply("ExtractUserScore", new ExtractAndSumScore("user"))  
  5. [... output ...]; 

ExtractAndSumScore实现了“What”中描述的逻辑,即按用户分组,然后累加分数,其相关代码如下:

 
  1. gameInfo  
  2. .apply(MapElements  
  3. .via((GameActionInfo gInfo) -> KV.of(gInfo.getKey(field), gInfo.getScore()))  
  4. .withOutputType(  
  5. TypeDescriptors.kvs(TypeDescriptors.strings(), TypeDescriptors.integers())))  
  6. .apply(Sum.integersPerKey()); 

通过MapElements确定Key与Value分别是用户与分数,然后Sum定义按key分组,并累加分数。Beam支持将多个对数据的操作合并成一个操作,这样不仅可以支持更清晰的业务逻辑实现,同时也可以在多处重用合并后的操作逻辑。

每小时团队分数

按照小时统计每个团队的分数,获得最高分数的团队可能获得奖励,这个分析任务增加了对窗口的要求,不过我们依然可以通过一个批处理任务实现,对于这个任务的“WWWH”四个维度的分析如下:

相对于第一个用户分数任务,只是在Where部分回答了“数据在什么范围中计算?”的问题,同时在What部分“如何计算数据?”中,分组的条件由用户改为了团队,这在代码中也会相应的体现:

 
  1. gameEvents  
  2. [... input ...]  
  3. [... parse ...]  
  4. .apply("AddEventTimestamps", WithTimestamps.of((GameActionInfo i)  
  5. -> new Instant(i.getTimestamp())))  
  6. .apply("FixedWindowsTeam", Window.into 
  7. FixedWindows.of(Duration.standardMinutes(windowDuration))))  
  8. .apply("ExtractTeamScore", new ExtractAndSumScore("team"))  
  9. [... output ...]; 

“AddEventTimestamps”定义了如何从原始数据中抽取EventTime数据,“FixedWindowsTeam”则定义了1小时固定窗口,然后重用了ExtractAndSumScore类,只是将分组的列从用户改成了团队。对于每小时团队分数任务,引入了关于“Where”部分窗口定义的新业务逻辑,但是从代码中可以看到,关于“Where”部分的实现和关于“What”部分的实现是完全独立的,用户只需要新加两行关于“Where”的代码,非常简单和清晰。

排行榜

前面两个任务均是基于有限数据集的批处理任务,对于排行榜来说,我们同样需要统计用户分数以及每小时团队分数,但是从业务角度希望得到的是实时数据。对于Apache Beam来说,一个相同处理逻辑的批处理任务和流处理任务的唯一不同就是任务的输入和输出,中间的业务逻辑Pipeline无需任何改变。对于当前示例的排行榜数据分析任务,我们不仅希望他们满足和前两个示例相同的业务逻辑,同时也可以满足更定制化的业务需求,例如:

  1. 流处理任务相对于批处理任务,一个非常重要的特性是,流处理任务可以更加实时的返回计算结果,例如计算每小时团队分数时,对于一小时的时间窗口,默认是在一小时的数据全部到达后,把最终的结算结果输出,但是流处理系统应该同时支持在一小时窗口只有部分数据到达时,就将部分计算结果输出,从而使得用户可以得到实时的分析结果。
  2. 保证和批处理任务一致的计算结果正确性。由于乱序数据的存在,对于某一个计算窗口,如何确定所有数据是否到达(Watermark)?迟到数据如何处理?处理结果如何输出,总量,增量,并列?流处理系统应该提供机制保证用户可以在满足低延迟性能的同时达到最终的计算结果正确性。

上述两个问题正是通过回答“When”和“How”两个问题来定义用户的数据分析需求。“When”取决于用户希望多常得到计算结果,在回答“When”的时候,基本上可以分为四个阶段:

  1. Early。在窗口结束前,确定何时输出中间状态数据。
  2. On-Time。在窗口结束时,输出窗口数据计算结果。由于乱序数据的存在,如何判断窗口结束可能是用户根据额外的知识预估的,且允许在用户设定的窗口结束后出现迟到的属于该窗口的数据。
  3. Late。在窗口结束后,有迟到的数据到达,在这个阶段,何时输出计算结果。
  4. Final。能够容忍迟到的最大限度,例如1小时。到达最后的等待时间后,输出最终的计算结果,同时不再接受之后的迟到数据,清理该窗口的状态数据。

对于每小时团队得分的流处理任务,本示例希望的业务逻辑为,基于Event Time的1小时时间窗口,按团队计算分数,在一小时窗口内,每5分钟输出一次当前的团队分数,对于迟到的数据,每10分钟输出一次当前的团队分数,在窗口结束2小时后迟到的数据一般不可能会出现,假如出现的话,直接抛弃。“WWWH”表达如下:

在基于Beam SDK的实现中,用户基于“WWWH” Beam Model表示的业务逻辑可以分别独立直接的实现出来:

 
  1. gameEvents  
  2. [... input ...]  
  3. .apply("LeaderboardTeamFixedWindows", Window  
  4. .into(FixedWindows.of 
  5. Duration.standardMinutes(Durations.minutes(60)))) 
  6. .triggering(AfterWatermark.pastEndOfWindow()  
  7. .withEarlyFirings(AfterProcessingTime.pastFirstElementInPane()  
  8. .plusDelayOf(Durations.minutes(5)))  
  9. .withLateFirings(AfterProcessingTime.pastFirstElementInPane()  
  10. .plusDelayOf(Durations.minutes(10))))  
  11. .withAllowedLateness(Duration.standardMinutes(120)  
  12. .accumulatingFiredPanes())  
  13. .apply("ExtractTeamScore", new ExtractAndSumScore("team"))  
  14. [... output ...] 

LeaderboardTeamFixedWindows对应“Where”定义窗口,Trigger对应“Where”定义结果输出条件,Accumulation对应“How”定义输出结果内容,ExtractTeamScore对应“What”定义计算逻辑。

总结

Apache Beam的Beam Model对无限乱序数据流的数据处理进行了非常优雅的抽象,“WWWH”四个维度对数据处理的描述,非常清晰与合理,Beam Model在统一了对无限数据流和有限数据集的处理模式的同时,也明确了对无限数据流的数据处理方式的编程范式,扩大了流处理系统可应用的业务范围,例如,Event-Time/Session窗口的支持,乱序数据的处理支持等。Apache Flink,Apache Spark Streaming等项目的API设计均越来越多的借鉴或参考了Apache Beam Model,且作为Beam Runner的实现,与Beam SDK的兼容度也越来越高。本文主要介绍了Beam Model,以及如何基于Beam Model设计现实中的数据处理任务,希望能够让读者对Apache Beam项目能够有一个初步的了解。由于Apache Beam已经进入Apache Incubator孵化,所以读者也可以通过官网或是邮件组了解更多Apache Beam的进展和状态。


本文作者:李呈祥

来源:51CTO

相关实践学习
简单用户画像分析
本场景主要介绍基于海量日志数据进行简单用户画像分析为背景,如何通过使用DataWorks完成数据采集 、加工数据、配置数据质量监控和数据可视化展现等任务。
SaaS 模式云数据仓库必修课
本课程由阿里云开发者社区和阿里云大数据团队共同出品,是SaaS模式云原生数据仓库领导者MaxCompute核心课程。本课程由阿里云资深产品和技术专家们从概念到方法,从场景到实践,体系化的将阿里巴巴飞天大数据平台10多年的经过验证的方法与实践深入浅出的讲给开发者们。帮助大数据开发者快速了解并掌握SaaS模式的云原生的数据仓库,助力开发者学习了解先进的技术栈,并能在实际业务中敏捷的进行大数据分析,赋能企业业务。 通过本课程可以了解SaaS模式云原生数据仓库领导者MaxCompute核心功能及典型适用场景,可应用MaxCompute实现数仓搭建,快速进行大数据分析。适合大数据工程师、大数据分析师 大量数据需要处理、存储和管理,需要搭建数据仓库?学它! 没有足够人员和经验来运维大数据平台,不想自建IDC买机器,需要免运维的大数据平台?会SQL就等于会大数据?学它! 想知道大数据用得对不对,想用更少的钱得到持续演进的数仓能力?获得极致弹性的计算资源和更好的性能,以及持续保护数据安全的生产环境?学它! 想要获得灵活的分析能力,快速洞察数据规律特征?想要兼得数据湖的灵活性与数据仓库的成长性?学它! 出品人:阿里云大数据产品及研发团队专家 产品 MaxCompute 官网 https://www.aliyun.com/product/odps 
相关文章
|
4月前
|
机器学习/深度学习 SQL 分布式计算
Apache Spark 的基本概念和在大数据分析中的应用
介绍 Apache Spark 的基本概念和在大数据分析中的应用
162 0
|
4月前
|
机器学习/深度学习 SQL 分布式计算
介绍 Apache Spark 的基本概念和在大数据分析中的应用。
介绍 Apache Spark 的基本概念和在大数据分析中的应用。
|
5天前
|
分布式计算 Java Go
Golang深入浅出之-Go语言中的分布式计算框架Apache Beam
【5月更文挑战第6天】Apache Beam是一个统一的编程模型,适用于批处理和流处理,主要支持Java和Python,但也提供实验性的Go SDK。Go SDK的基本概念包括`PTransform`、`PCollection`和`Pipeline`。在使用中,需注意类型转换、窗口和触发器配置、资源管理和错误处理。尽管Go SDK文档有限,生态系统尚不成熟,且性能可能不高,但它仍为分布式计算提供了可移植的解决方案。通过理解和掌握Beam模型,开发者能编写高效的数据处理程序。
134 1
|
27天前
|
数据可视化 Linux Apache
CentOS部署Apache Superset大数据可视化BI分析工具并实现无公网IP远程访问
CentOS部署Apache Superset大数据可视化BI分析工具并实现无公网IP远程访问
|
1月前
|
机器学习/深度学习 分布式计算 大数据
一文读懂Apache Beam:统一的大数据处理模型与工具
【4月更文挑战第8天】Apache Beam是开源的统一大数据处理模型,提供抽象化编程模型,支持批处理和流处理。它提倡"一次编写,到处运行",可在多种引擎(如Spark、Dataflow、Flink)上运行。Beam的核心特性包括抽象化概念(PCollection、PTransform和PipelineRunner)、灵活性(支持多种数据源和转换)和高效执行。它广泛应用在ETL、实时流处理、机器学习和大数据仓库场景,助力开发者轻松应对数据处理挑战。
25 1
|
1月前
|
分布式计算 资源调度 Hadoop
Apache Hadoop入门指南:搭建分布式大数据处理平台
【4月更文挑战第6天】本文介绍了Apache Hadoop在大数据处理中的关键作用,并引导初学者了解Hadoop的基本概念、核心组件(HDFS、YARN、MapReduce)及如何搭建分布式环境。通过配置Hadoop、格式化HDFS、启动服务和验证环境,学习者可掌握基本操作。此外,文章还提及了开发MapReduce程序、学习Hadoop生态系统和性能调优的重要性,旨在为读者提供Hadoop入门指导,助其踏入大数据处理的旅程。
186 0
|
2月前
|
分布式计算 大数据 Apache
大数据技术变革正当时,Apache Hudi了解下?
大数据技术变革正当时,Apache Hudi了解下?
25 0
|
2月前
|
存储 数据处理 Apache
万字长文 | 泰康人寿基于 Apache Hudi 构建湖仓一体平台的应用实践
万字长文 | 泰康人寿基于 Apache Hudi 构建湖仓一体平台的应用实践
98 0
|
3月前
|
SQL 并行计算 大数据
【大数据技术攻关专题】「Apache-Flink零基础入门」手把手+零基础带你玩转大数据流式处理引擎Flink(基础加强+运行原理)
关于Flink服务的搭建与部署,由于其涉及诸多实战操作而理论部分相对较少,小编打算采用一个独立的版本和环境来进行详尽的实战讲解。考虑到文字描述可能无法充分展现操作的细节和流程,我们决定以视频的形式进行分析和介绍。因此,在本文中,我们将暂时不涉及具体的搭建和部署步骤。
498 3
【大数据技术攻关专题】「Apache-Flink零基础入门」手把手+零基础带你玩转大数据流式处理引擎Flink(基础加强+运行原理)
|
4月前
|
SQL 关系型数据库 Apache
Apache Doris 整合 FLINK CDC 、Paimon 构建实时湖仓一体的联邦查询入门
Apache Doris 整合 FLINK CDC 、Paimon 构建实时湖仓一体的联邦查询入门
693 1

推荐镜像

更多