ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Apache Beam 0.4.0 新增 Apache Apex Runner:统一模型到状态流引擎的翻译之路

Apache Beam 0.4.0 新增 Apache Apex Runner:统一模型到状态流引擎的翻译之路 【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载Apache Beam 在 0.4.0 版本中正式加入了对 Apache Apex 的 runner 支持这是 Beam 可移植统一编程模型在流处理领域的一次重要落地。本文以仓库内历史博客 added-apex-runner.md 为核心系统讲解该 runner 诞生的背景、Apex 作为状态流处理器的底层原理、Beam 模型翻译为 Apex DAG 的机制、嵌入式执行与测试方式并结合当前仓库证据梳理这一 runner 从加入到移除的完整历史脉络。读完本文你将理解一个 Beam runner 如何将统一模型映射到具体执行引擎以及 Apex 的容错状态机制如何支撑 Beam 的窗口与事件时间语义。背景Beam 的统一编程模型与 Apex 流处理框架Apache Beam 起源于 Google Dataflow SDK作为 Apache 孵化器项目快速接纳了 Apache 社区运作方式并围绕统一编程模型 多引擎可移植这一核心理念吸引了大量用户——同一份 Beam 管道代码可以运行在批处理与流处理等不同大数据框架之上这正是其价值所在即业界常说的 Streaming-101 / Streaming-102 所阐述的超越批处理思想。当时多个 Apache 项目已经为 Beam 提供了 runner 实现参见 runners 能力矩阵。Apache Apex 则是一个面向集群的流处理框架专为低延迟、高吞吐、有状态、可靠地对复杂分析管道进行处理而设计。Apex 自 2012 年开始开发已被大型公司在实时与批量处理场景中用于生产环境。正是这样的定位使它成为 Beam 流式语义落地时一个有吸引力的执行目标。0.4.0 版本新增 Apex runner 是 Beam 社区的一个重要里程碑。需要说明的是本文描述的是该 runner 加入时的初始状态初版实现聚焦于在功能层面广泛覆盖 Beam 模型后续仍需多个方向的工作才能从功能可用走向可扩展、高性能以匹配 Apex 及其原生 API 的能力。Apex 作为状态流处理器的核心能力Apex 从设计之初就是一个状态流处理器stateful stream processor其关键机制构成了整个 runner 容错与一致性能力的基础分布式异步检查点checkpoint操作符operator以分布式、异步的方式对状态做检查点从而为整个处理图DAG产出一份一致的快照该快照可用于故障恢复。增量/细粒度恢复Apex 支持增量式的恢复方式——故障发生时只有 DAG 中实际受影响的局部会被恢复其余管道继续处理。这种特性可以支撑特殊需求的用例例如通过推测执行speculative execution来满足处理延迟上的 SLA。幂等处理保证状态检查点与幂等处理相结合是 Apex 支持 exactly-once恰好一次结果的基础。这一能力对 Beam runner 的意义在于要广泛支持 Beam 的窗口windowing概念、尤其是基于事件时间event time的处理就必须能够以容错且高效的方式跟踪计算状态。从能力矩阵看Apex 的能力与 Beam 模型对齐得相当好。翻译到 Apex DAGrunner 的职责与最小操作符集合一个 Beam runner 的核心职责是实现从 Beam 模型到底层框架执行模型execution model的翻译。对 Apex runner 而言翻译目标就是 Apex 原生的、可组合的低层 DAG API——它同时也是 Apex 上用于指定应用的其他多种 API 的基础。Apex 的 DAG 由操作符functional building blocks功能构建块构成操作符之间以流streams连接runner 则提供执行层。在 Apex 场景下执行层就是分布式流处理操作符逐事件event by event处理数据。初版 runner 覆盖 Beam 原始变换primitive transforms的最小操作符集合包括ParDo.Bound有界输入的 ParDo对每个元素执行用户定义的 DoFn。ParDo.BoundMulti多输出旁路输出的 ParDo。Read.Unbounded无界数据源读取对应流式输入。Read.Bounded有界数据源读取对应批式输入。GroupByKey按键分组是窗口与聚合语义的核心变换。Flatten.FlattenPCollectionList将多个 PCollection 合并为一个。从 Beam 的模型结构看用户组合出的各类高层变换最终都会归结为这些原始变换因此覆盖它们即意味着具备了对通用 Beam 管道的基础执行能力。执行与测试嵌入式模式与集成测试套件在本版本中Apex runner 以**嵌入式模式embedded mode**执行管道与直接执行器Direct Runner类似所有内容在单个 JVM 内运行。这种模式非常适合开发与调试阶段的快速验证。如何用 Apex runner 运行 Beam 示例可参考 Java Quickstart当前仓库中已无 Apex 专属指南以通用 quickstart 为准。嵌入式模式之外的两种形态需要说明生产环境Apex 在生产环境中分布式运行于 Apache Hadoop YARN 集群之上。初版曾给出一个将 Beam 管道嵌入 Apex 应用包以运行于 YARN 的示例apex-samples 仓库中的 beam-apex-wordcount而 runner 内的直接启动direct launch支持当时仍在开发中。测试保障Beam 项目高度重视开发流程与工具链包括测试。针对各 runner 有一套综合测试套件包含200 多个集成测试每次变更都会针对每个 runner 执行确保功能不被破坏。这些测试覆盖了能力矩阵中的各项能力因而是衡量 runner 实现完整性与正确性的标尺。该套件在 Apex runner 的开发过程中发挥了很大作用。展望与后续演进从功能可用到生产就绪初版 Apex runner 的下一步规划是将其从功能可用推进到可支撑分布式真实应用充分利用 Apex 的扩展性与性能特性与其原生 API 相当。这包括ParDo 的链接/合并chaining of ParDos分区partitioningcombine 操作的优化从当前仓库证据看Apex runner 的历史轨迹可以完整还原beam-a-look-back.md 记录了 2017 年初 Beam 生态的 runner 阵容Apache Flink、Apache Spark 1.x、Google Cloud Dataflow、Apache Apex、Apache Gearpump印证了 Apex 是 Beam 早期多引擎可移植策略的重要成员。beam-2.23.0.md 的 Deprecations 一节明确记载2.23.0 版本移除了 Apex runnerBEAM-9999与 Gearpump runner 一同退出。从源码结构看当前仓库 runners 目录下已无 apex 子模块runner 文档列表中也仅保留 dataflow、direct、flink、jet、jstorm、mapreduce、nemo、samza、spark、twister2 等条目。因此本文所描述的技术细节反映的是该 runner 在 0.4.0 引入时的设计形态对于今天在 Apache Beam 中选择 runner 的开发者应以其时当前版本的 runner 能力矩阵 为准。Apex runner 的历史价值在于它验证了 Beam 统一模型翻译到有状态流引擎的可行性尤其是事件时间窗口语义对底层状态跟踪能力的需求为后来 runner 的容错与一致性设计提供了参考范本。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐Apache Beam 0.4.0 新增 Apache Apex Runner从统一编程模型到有状态流处理引擎的落地实践Apache Beam 0.4.0 新增 Apache Apex Runner从统一编程模型到有状态流处理引擎的落地实践 本篇技术指南基于 Apache Be大数据批处理流处理数据工程Apache Beam 0.4.0 引入 Apex Runner从 Beam 模型到 Apex DAG 的转换与内嵌执行Apache Beam 0.4.0 引入 Apex Runner从 Beam 模型到 Apex DAG 的转换与内嵌执行 2017 年 1 月发布的 Apac批处理流处理大数据Apache Beam 仓库全景指南统一批流编程模型、SDK 与 Runner 生态Apache Beam 仓库全景指南统一批流编程模型、SDK 与 Runner 生态 Apache Beam 是一个用于定义 批处理Batch与流处理S批处理流处理大数据上一篇如何快速掌握大麦网自动抢票神器3倍成功率实战指南下一篇告别手速焦虑大麦网自动抢票工具终极指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进