【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
导读:本文基于 Apache Beam 官方博客《Apache Beam: Six Months in Incubation》撰写,聚焦 Beam 进入 Apache 软件基金会(ASF)孵化器后第一个半年(2016 年上半年)的社区与工程进展。你将看到当时 Beam 起步阶段的关键数字、runner 无关化重构、Flink/Spark 等 runner 的能力补全、多 SDK 与数据源/汇(IO)生态的早期布局,以及社区建设策略——并对照当前仓库(gh_mirrors/beam18/beam)中可验证的源码结构,理解这些早期决策如何沉淀为今天的 Apache Beam 形态。
2016 年 8 月,Apache Beam 正式进入 Apache 软件基金会孵化(Incubation)刚满半年。作为后来被广泛用于批流统一数据处理的开源项目,这六个月的起点值得回顾:它既是"从三家公司的代码捐赠走向开放社区"的社区叙事,也是"从绑定特定执行引擎的代码库走向 runner 无关框架"的技术分水岭。本文以这篇半年回顾为骨架,结合当前仓库中保留的实现痕迹,还原那段关键时期的真实进展。
一、半年盘点:一组数字背后的起步
原文开篇给出了孵化前六个月的一组硬指标,这些数字构成了理解 Beam 早期阶段的第一手资料:
| 指标 | 数值 | 解读 |
|---|---|---|
| 捐赠代码量 | 48,238 行 | 来自 Cloudera、dataArtisans、Google 三家公司 |
| Pull Request | 761 个,来自 45 位贡献者 | 社区协作已经成型 |
| Jira 问题 | 498 个已开启,245 个已解决 | 问题追踪活跃、解决率约一半 |
| 版本发布 | 1 个孵化版(另有 1 个在途) | 首个孵化版已落地 |
| 自动化测试 | 4,200 小时 | 持续集成投入巨大 |
| user@ 邮件列表 | 161 订阅者 / 606 条消息 | 用户社区开始聚集 |
| dev@ 邮件列表 | 217 订阅者 / 1,205 条消息 | 开发者讨论密度远高于用户列表 |
| GitHub 关注/复刻 | 277 stars / 174 forks | 早期影响力正在积累 |
这些数字透露了一个重要信息:Beam 的起点并非"从零写起",而是三家在流处理领域各有积累的公司(Cloudera 的 Crunch、dataArtisans 的 Flink 系、Google 的 Dataflow 系)把既有代码并入 Apache 孵化器,再由开放社区接手重构。这与今天仓库中sdks/java、runners等目录的模块化布局一脉相承。
二、技术主线:runner 无关化与核心能力补全
2.1 整个代码库的 runner 无关化重构
半年回顾强调的第一项技术工作,是对"整个代码库、示例和测试进行真正的 runner 无关化重构"。这意味着:
- 核心编程模型(PCollection、PTransform、DoFn、窗口/触发器等)不得绑定任何具体执行引擎;
- 示例代码与测试必须在任何 runner 上都能运行;
- 各 runner 通过统一的执行语义与核心层对接。
对照当前仓库,这一决策的直接体现是runners/目录下并列维护着 core-construction-java、core-java、direct-java 以及 Flink、Spark、Samza、Jet、Twister2、Google Cloud Dataflow 等众多 runner 实现,每个 runner 都是"核心模型 + 引擎适配"的独立模块。当时的重构正是为今天这种"同一份作业可迁移到多种引擎"的架构奠定基础。
2.2 Flink runner:时间戳、窗口与有界源
原文特别点名了 Apache Flink runner 的两项新能力:
- 批处理与有界源(bounded sources)中的时间戳/窗口支持——使 Flink runner 在批模式下也能正确处理事件时间与窗口聚合;
- 流式模式下的 side inputs(侧输入)——让流式作业可以引用外部数据集(如配置表、查找表)。
从源码结构看,这些能力的实现沉淀在 runners/flink 模块(含src/main128 个 Java 文件、57 个测试文件,并针对 Flink 1.12–1.16 维护多版本适配),窗口、水印、侧输入等语义属于 core-java 定义的核心模型的一部分,runner 侧负责翻译与执行。
2.3 Spark runner:向 Spark 2.0 升级
同期正在推进的工作是将 Spark runner 升级到 Spark 2.0。当前仓库中 runners/spark 既保留 Spark 3 版本适配(3/子目录、73 个 Java 文件),也包含主实现与spark_runner.gradle构建定义,说明这条 runner 适配路线从孵化期延续至今,始终跟随 Spark 主线版本演进。
2.4 来自更广泛 Apache 社区的 runner 候选
孵化半年时,多个 Apache 项目已表现出接入 Beam 的意向:
- Apache Gearpump:建立了独立 feature branch;
- Apache Apex:已有 PR 提交;
- Apache Storm:开始初步讨论。
如今这些意向多数已落地为正式 runner(如仓库中runners/下可见的 Samza、Jet、Twister2 等模块),印证了"统一编程模型 + 多种执行引擎"的设计在当时就具有很强的社区吸引力。
三、SDK/DSL 生态:Python SDK 与 Scio DSL
原文记录的 SDK 进展包括:
- Python SDK(来自 Google):已进入 feature branch,是 Beam 暴露编程模型的第二门语言;
- Scio DSL(来自 Spotify):规划中,基于 Scala 的 Beam 封装。
对照当前仓库,sdks/python/apache_beam 已是拥有 1,100+ Python 文件的成熟 SDK(含 beam 核心、transforms、io、runner 实现与类型标注等),sdks/go 也成长为独立 SDK。孵化期"多语言暴露统一模型"的规划,最终演化为今天 Java / Python / Go / TypeScript / Scala(Scio)多 SDK 并存的局面。
四、数据源与汇(IO)扩展:Kafka、JMS 与更多连接器
半年回顾列出了 IO 生态的早期进展:
- 已完成:Apache Kafka、JMS;
- PR 进行中:Amazon Kinesis、Apache Cassandra、MongoDB;
- 规划中:更多连接器。
这两条在今天的仓库中都能找到直接对应:
- sdks/java/io/kafka:Kafka IO 模块,含
build.gradle构建定义; - sdks/java/io/jms:JMS IO 模块,同样独立成目录;
- sdks/java/io/kinesis、sdks/java/io/cassandra、sdks/java/io/mongodb 均已正式落地。
此外sdks/java/io下还扩展出 amazon-web-services、azure、debezium、elasticsearch、hbase、jdbc、parquet、pulsar、rabbitmq、redis、snowflake、splunk 等 40+ 连接器目录。当年"连接器正在被社区不断补充"的预期,如今已兑现为 Beam 的 IO 矩阵。值得一提的是,本仓库 it/kafka 等目录中还保留着连接器的集成测试,用于验证与外部系统交互的正确性。
五、社区建设:从设计讨论到全球宣讲
半年回顾用较多篇幅强调"构建一个参与度高、包容的社区",具体动作包括:
- 开发者社区:围绕 DoFn 复用语义(DoFn reuse semantics)、序列化技术、状态访问 API(state API)等展开详细设计讨论——这些主题最终都固化为核心编程模型的一部分;
- 用户社区:运营活跃邮件列表,持续改进官网与文档;
- 技术传播:在 ApacheCon、Hadoop Summit、Kafka Summit、JBCN Barcelona、Strata 等多个会议发表 Beam 主题演讲;
- 线下活动:参加多个既有 meetup,并开始组织自己的 meetup。
从今天的仓库看,社区治理沉淀为 website/www/site/content/en/contribute 下的一整套贡献指南(含 get-started-contributing.md、ptransform-style-guide.md、runner-guide.md 等),以及仓库根目录的 CONTRIBUTING.md。这些文档系统正是"欢迎新贡献者、明确协作方式"这一社区目标的制度化产物。
六、面向稳定版与毕业:如何参与
半年回顾的结尾提出了两个明确目标:发布稳定版本与从孵化器毕业,并向读者发出参与邀请:
- 加入邮件列表(user@ 与 dev@);
- 阅读贡献指南;
- 从 Jira 上的 starter task(标记为 newbie/starter)入手。
这一目标在次年得到了兑现:2017 年 1 月,Beam 正式从孵化器毕业,成为 Apache 顶级项目(详见仓库中 beam-graduates.md 的公告博客),并随后发布了 2.0 稳定版。对今天的潜在贡献者而言,仓库中的 CONTRIBUTING.md 仍是最直接的入口,它完整描述了从"分享意图(先建 issue 并认领)"、搭建开发环境(Java/Go/Python/Docker,可用 local-env-setup.sh 一键配置)到"提交 PR 并触发 pre-commit 测试"的全流程。
七、结语:半年纪事与今天的 Beam
回望这份六个月纪事,可以提炼出几条至今仍成立的 Beam 发展主线:
- 统一模型先行:runner 无关化重构不是"优化",而是"地基",它让多引擎、多语言、多连接器的生态扩张成为可能;
- 社区驱动扩展:Gearpump/Apex/Storm 的候选 runner、Kinesis/Cassandra/MongoDB 的连接器 PR,无一不是社区成员主动贡献的结果;
- 以测试与发布节奏保障质量:4,200 小时自动化测试和半年内首个孵化版,奠定了后来"每 6 周一个 minor release"的工程节奏。
六年多过去,这份半年回顾中的"规划"大多已成现实:今天的 Beam 支持 Java/Python/Go/TypeScript 等多语言 SDK,拥有 Flink/Spark/Dataflow/Direct 等主流 runner,以及覆盖 Kafka、Kinesis、Cassandra、MongoDB、JDBC、Elasticsearch 等数十种数据源的 IO 生态。如果你正在评估或使用 Beam,这份早期纪事能帮你理解其架构哲学与社区文化的来源——而当前仓库本身,就是这段历史最完整的注脚。
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam 孵化六个月回顾:从 48,238 行捐赠代码到统一批流编程模型
Apache Beam 孵化六个月回顾:从 48,238 行捐赠代码到统一批流编程模型 导读:本文基于 Apache Beam 社区发布的《Apache Bea
大数据批处理流处理数据工程背完就忘?这套408思维导图笔记CS-Xmind-Note,把4本教材压成几十张图
背完就忘?这套408思维导图笔记CS Xmind Note,把4本教材压成几十张图 深夜十点,图书馆的灯还亮着。你把《计算机操作系统》翻到"进程调度"那一章,读
文档教程知识库SkyDNS版本升级指南:从v1到v2的重大变更解析
SkyDNS版本升级指南:从v1到v2的重大变更解析 SkyDNS是一款轻量级DNS服务器,从v1到v2版本的升级带来了多项重要改进与架构调整。本文将详细解析版
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考