【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
导读
本文面向希望以 JavaScript/TypeScript 原生方式参与 Apache Beam 数据处理的开发者,系统讲解 TypeScript SDK 的设计目标、从源码构建与运行的完整工作流,以及它如何借助 Beam 可移植性框架(portability framework)在 Direct、Flink、Dataflow 等多种运行器上执行 Pipeline。读完本文,你将掌握npm run build构建、node ... --runner=...运行 wordcount 示例、npm test跑测试的具体操作,并理解 DirectRunner 复用 Worker、跨语言 Transform(cross-language transforms)经 Expansion Service 展开、可移植运行器经 Job Service 提交作业等底层机制。
一、TypeScript SDK 的双重使命
Apache Beam 是一个面向批处理与流处理的统一编程模型,官方 SDK 覆盖 Java、Python 和 Go。而sdks/typescript目录下的 TypeScript SDK(sdks/typescript/README-dev.md)是 Beam 团队对 JavaScript 生态的一次原生尝试,它承载着两个截然不同的目标:
- 触达庞大的 JavaScript 开发者社区。现有数据处理框架对 JavaScript 开发者的支持相对不足,一个原生针对该语言的 SDK 可以填补这一空白。
- 充当 Beam 移植到新语言的参考实现(proof of concept)。Beam 与 Dataflow 的一个重要特性是易于移植到新语言,这个 SDK 本身就是展示这种可移植性的活的样例。
为了实现上述目标,SDK 在架构上高度依赖可移植性框架(portability framework),其核心思路是:Pipeline 的定义与执行彻底分离,先构建出与语言无关的 Beam Runner API 协议(protobuf)描述,再交给任意支持该协议的后端运行器执行。
从源码结构看,这一设计贯穿始终:
sdks/typescript/src/apache_beam/proto/存放由 protobuf 定义生成的全部协议代码(beam_runner_api、beam_fn_api、beam_job_api、beam_expansion_api、beam_artifact_api等),由 gen_protos.sh 生成;sdks/typescript/src/apache_beam/runners/存放各类运行器的实现;sdks/typescript/src/apache_beam/worker/存放可复用的 Worker 执行逻辑。
两个关键设计印证了“以可移植性为中心”的定位:
- IO 大量使用跨语言 Transform:TypeScript SDK 本身只实现少量本地 IO,更多的读写能力(如 Kafka、BigQuery、Pub/Sub 等)通过跨语言 Transform 委托给其他 SDK 的 Expansion Service 展开;
- Direct Runner 就是 Worker 的扩展:本地直接运行并非另起炉灶,而是直接复用 SDK Worker 中
BundleProcessor的能力(见下文第四节),因此本地能跑的 Pipeline 可以平滑迁移到 Dataflow、Flink 等生产运行器。
对使用者而言,这意味着运行其他语言代码被封装在 Docker 镜像中并非障碍——这正是该 SDK 有意选择的设计路线。
二、从源码构建与安装
2.1 前置环境
TypeScript SDK 的本地开发需要npm与python两类工具。其中Python 并非可选依赖:它被用来编排 Beam 功能(orchestrate Beam functionality),例如在可移植运行器模式下启动 Job Service。
注意:
README-dev.md中给出的 clone 命令为git checkout https://github.com/apache/beam(原文笔误,实际应为git clone),当前仓库即为 Beam 源码,以下操作均以已获得源码为前提。
在仓库根目录下进入 TypeScript SDK 目录,安装 npm 依赖:
cd sdks/typescript npm installnpm install会根据 sdks/typescript/package.json 拉取运行时依赖(如@grpc/grpc-js、@protobuf-ts/grpc-transport、bson、protobufjs、serialize-closures、ttypescript等)与开发依赖(typescript@4.7、mocha、eslint、prettier、typedoc等)。
2.2 构建:TypeScript 编译为 JavaScript
npm run buildpackage.json中build脚本实际执行bash build.sh,将src/下的 TypeScript 文件转译(transpile)为 JS 并输出到dist/目录。编译配置见 sdks/typescript/tsconfig.json,其中有几个值得注意的点:
module: commonjs、target: es2021、outDir: "dist/";- 启用了
ts-closure-transform编译器插件(beforeTransform与afterTransform两个阶段),用于将闭包(函数/生成器)序列化为可在分布式 Worker 间传递的形式——这是 JS 函数能够跨进程执行的关键; declaration: true与sourceMap: true会同时产出.d.ts类型声明与源码映射。
构建完成后,dist/src/apache_beam/index.js即成为 npm 包的入口(package.json的main字段),apache-beam-worker可执行入口指向dist/src/apache_beam/worker/worker_main.js,用于启动独立 Worker 进程。
2.3 通过 npm 安装使用
对于普通使用者,无需从源码构建,直接安装发布包即可:
npm install apache_beam由于 SDK 大量使用跨语言 Transform,官方建议系统上同时具备Python 3 与 Java,以便在需要时启动 Expansion Service 或 Job Service。
2.4 开发工作流
所有开发工作流(build、test、lint、clean 等)都定义在package.json的scripts字段中,可通过 npm 命令调用:
| npm 命令 | 实际行为 |
|---|---|
npm run build | 执行bash build.sh,编译 TS → JS 到dist/ |
npm run clean | tsc --clean,清理编译产物 |
npm test | 先pretest自动构建,再运行mocha dist/test dist/test/docs执行测试 |
npm run lint | 运行eslint . --ext .ts检查代码 |
npm run prettier | 用prettier --write src/自动格式化源码 |
npm run prettier-check | 用prettier --check src/校验格式 |
npm run docs | 先构建,再用typedoc生成 API 文档 |
npm run codecovTest | 通过 istanbul 生成覆盖率并上传 codecov |
npm run worker | 直接以node运行外部 Worker 服务external_worker_service.js |
三、运行一个真实 Pipeline:wordcount
sdks/typescript/src/apache_beam/examples/wordcount.ts定义了一个参数化的 wordcount Pipeline,可以通过--runner参数在不同的运行器上执行。构建完成后,直接运行编译产物即可:
node dist/src/apache_beam/examples/wordcount.js ${PARAMETERS}3.1 在本地 Direct Runner 上运行
node dist/src/apache_beam/examples/wordcount.js --runner=direct--runner=direct对应 sdks/typescript/src/apache_beam/runners/direct_runner.ts 中的directRunner工厂函数。它不依赖任何外部服务,适合快速验证 Pipeline 逻辑。
3.2 在 Flink 上运行(基础设施自动下载)
node dist/src/apache_beam/examples/wordcount.js --runner=flink--runner=flink会走 sdks/typescript/src/apache_beam/runners/flink.ts 中的flinkRunner:本地基础设施(Flink Job Server)会被自动下载并启动,无需手工搭建。其默认参数为flinkMaster: "[local]"与flinkVersion(取runners/flink下已发布的版本列表如 1.12/1.13/1.14 中的最新版本,见源码中与gradle.properties保持同步的PUBLISHED_FLINK_VERSIONS常量)。Job Server 以 Java jar 形式(通过JavaJarService从runners:flink:${flinkVersion}:job-server:shadowJar构建/缓存)拉起,并以--flink-master、--artifacts-dir、--job-port、--artifact-port等参数启动。
3.3 在 Google Cloud Dataflow 上运行
node dist/src/apache_beam/examples/wordcount.js \ --runner=dataflow \ --project=${PROJECT_ID} \ --tempLocation=gs://${GCS_BUCKET}/wordcount-js/temp --region=${REGION}--runner=dataflow走 sdks/typescript/src/apache_beam/runners/dataflow.ts 中的dataflowRunner,需要提供project、tempLocation、region三个必选参数。从源码可以看到,Dataflow 模式通过PythonService.forModule("apache_beam.runners.dataflow.dataflow_job_service", ...)启动 Python 实现的 Job Service,并在 Pipeline 选项中自动注入三个实验开关:
use_runner_v2:启用 Runner V2 执行框架;use_portable_job_submission:使用可移植作业提交方式;use_sibling_sdk_workers:使用兄弟 SDK Worker。
这再次印证了“生产运行器 = 可移植框架 + Job Service”的统一路径。
3.4 wordcount Pipeline 本身长什么样
wordcount.ts的主体代码非常简洁(源码):
function wordCount(lines: beam.PCollection<string>): beam.PCollection<any> { return lines .map((s: string) => s.toLowerCase()) .flatMap(function* (line: string) { yield* line.split(/[^a-z]+/); }) .apply(countPerElement()); } async function main() { await createRunner(yargs.argv).run((root) => { const lines = root.apply( beam.create([ "In the beginning God created the heaven and the earth.", "And the earth was without form, and void; and darkness was upon the face of the deep.", // ... ]), ); lines.apply(wordCount).map(console.log); }); }这段代码展示了 TypeScript SDK 的几个 API 特色:
root.apply(beam.create([...]))创建包含原始文本的 PCollection;.map(...)直接对元素做逐条变换(转小写);.flatMap(function* (line) { yield* ... })以生成器(generator)方式产出多个元素——与 Python SDK 的flatMap/ParDo.process风格一致,而不是回调式;.apply(countPerElement())应用组合类 Transform 完成词频统计;createRunner(yargs.argv)根据命令行参数选择运行器,然后run(...)等待作业彻底完成。
四、运行器体系:从 Direct Runner 到可移植 Runner
4.1 运行器的创建与选择
所有运行器都通过 sdks/typescript/src/apache_beam/runners/runner.ts 中的createRunner(options)统一创建,支持以下取值:
--runner值 | 实现 |
|---|---|
default(缺省) | defaultRunner:能在 Direct 上跑的优先用 Direct,否则回退到 universal(可移植)运行器 |
direct | DirectRunner(direct_runner.ts) |
universal | universalRunner:通过 Python 启动本地 Job Service 的可移植运行器(universal.ts) |
flink | flinkRunner:自动拉起 Flink Job Server(flink.ts) |
dataflow | dataflowRunner:对接 Google Cloud Dataflow(dataflow.ts) |
Runner抽象类提供了两种执行入口:
run(pipelineFn, options):等 Pipeline 完全结束后才返回(最终状态非DONE会抛出异常),不易出错,适合脚本场景;runAsync(pipelineFn, options):立即返回一个PipelineResult句柄,用于查询作业状态与指标。
PipelineResult还封装了waitUntilFinish(duration)(毫秒超时轮询)、counters()、distributions()等指标聚合方法(基于beam:metric:user:sum_int64:v1/beam:metric:user:distribution_int64:v1两类 MonitoringInfo 聚合)。
4.2 defaultRunner 的智能回退
defaultRunner的实现体现了“Direct 优先、能力不足时升级”的设计:
const directRunner = require("./direct_runner").directRunner(defaultOptions); if (directRunner.unsupportedFeatures(pipeline, options).length === 0) { return directRunner.runPipeline(pipeline, options); } else { return loopbackRunner(defaultOptions).runPipeline(pipeline, options); }DirectRunner.unsupportedFeatures会检查 Pipeline 中是否包含 Direct Runner 不支持的要素,例如:
requirements中存在未登记的能力要求(SUPPORTED_REQUIREMENTS为空数组,意味着任何额外需求都会触发回退);- 环境中出现非
TYPESCRIPT_DEFAULT_ENVIRONMENT_URN的 URN(说明涉及跨语言执行); - 窗口合并策略(
MergeStatus)或输出时间(OutputTime)配置超出 Direct 支持范围。
一旦发现不支持的要素,Pipeline 自动交给universalRunner(environmentType: "LOOPBACK"的 loopback 模式),从而保证“同一份 Pipeline 代码本地可跑、生产可迁”。
4.3 Direct Runner 本质上是 Worker 的扩展
README-dev.md明确指出:"the direct runner is simply an extension of the worker suitable for running on portable runners such as the ULR"。这一点在DirectRunner.runPipeline的源码中得到印证:
const processor = new worker.BundleProcessor( descriptor, null!, new state.CachingStateProvider(stateProvider), [impulse.urn], ); await processor.process("bundle_id");Direct Runner 直接把 Pipeline 的 Runner API protobuf 组装成ProcessBundleDescriptor,交给 sdks/typescript/src/apache_beam/worker/worker.ts 中的BundleProcessor作为单个 bundle 处理。也就是说,本地直跑与分布式 Worker 执行共用同一套算子(operator)执行引擎,这是它能平滑迁移到可移植运行器的根本原因。
为了支持单个 bundle 内的执行语义,direct_runner.ts还实现了几个专用算子:
DirectImpulseOperator:模拟 Beam 的 impulse 源(每个 pipeline 触发一个元素,而不是每个 worker);DirectGbkOperator:在单 bundle 内完成 GroupByKey(对 key 按 window+key 分组);rewriteSideInputs:为含旁路输入(side input)的 ParDo 重写执行图——插入CollectSideOperator收集旁路输入到内存状态、插入BufferOperator缓冲主输入,确保旁路输入收集完毕后 parDo 才执行;InMemoryStateProvider:以内存 Map 提供状态读写,支撑 side input 与未来状态功能。
这也解释了README-dev.md中 TODO 列表里"真正使用 worker threads 并行处理多个 bundle"的动机:当前 Direct Runner 把整个 Pipeline 当作一个 bundle 顺序执行,并行化是后续演进方向。
4.4 可移植 Runner:经过 Job Service 的远程执行
对于universal/flink/dataflow,执行路径最终收敛到 sdks/typescript/src/apache_beam/runners/portable_runner/runner.ts 中的PortableRunner。其核心流程(runPipelineWithProto)可归纳为:
- 探测流式需求:扫描所有 PCollection,若存在
isBounded == UNBOUNDED,自动把streaming选项置为 true; - 选择执行环境:
LOOPBACK模式:通过ExternalWorkerPool(external_worker_service.ts 中的 gRPC ExternalWorkerPool 服务)在本地进程内启动 Worker,适合本地验证;- 默认模式:将 SDK 环境替换为 Docker 环境,镜像为
docker.io/apache/beam_typescript_sdk:<版本>(可用sdkContainerImage选项覆盖),并自动执行npm pack把当前代码打成 npm 包作为 artifact 注册,同时收集file:形式的本地依赖一并上传;
- 注册模块:把需要 Worker 端 import 的模块集合写入
registeredNodeModules(来自serialization.getRegisteredModules()); - Prepare + Run:调用 Job Service 的
Prepare方法提交PrepareJobRequest(Pipeline 与转成beam:option:*前缀的 pipeline options);若响应要求暂存 artifact,则通过 Artifact Staging Service 上传;最后调用Run方法拿到jobId; - 返回句柄:
runPipelineWithProto在作业成功提交后即返回PortableRunnerPipelineResult,由用户通过waitUntilFinish轮询getState,直到进入DONE/FAILED/CANCELLED/UPDATED/DRAINED等终态。
Job Service 本身由各运行器负责拉起:
universalRunner用PythonService.forModule("apache_beam.runners.portability.local_job_service_main", ...)启动本地 Python Job Service(ULR,Universal Local Runner);flinkRunner用JavaJarService启动 Flink Job Server jar;dataflowRunner用PythonService.forModule("apache_beam.runners.dataflow.dataflow_job_service", ...)启动 Dataflow Job Service。
4.5 跨语言 Transform:IO 的主力实现方式
由于 TypeScript SDK 将 IO 大量委托给跨语言 Transform,理解 Expansion Service 机制是掌握该 SDK 的关键。以 sdks/typescript/src/apache_beam/transforms/external.ts 中的rawExternalTransform为例,跨语言 Transform 的展开过程为:
- 构造
ExpansionRequest,携带要展开的PTransform(含urn与 payload 配置)及其输入 PCollection; - 调用指定地址的 Expansion Service(gRPC)的
expand方法; - 若展开结果的环境带有依赖(
dependencies),则通过 Artifact Retrieval Service 拉取这些 artifact 转存为持久形式(resolveArtifacts); - 将返回的 Pipeline 片段(transforms、pcollections、coders、environments、windowingStrategies)按 namespace 拼接到当前 Pipeline 中(
splice),并校验输出 coders 可被 SDK 理解。
sdks/typescript/src/apache_beam/io/index.ts中导出的 IO 清单(avroio、bigqueryio、kafka、parquetio、pubsub、pubsublite、schemaio、textio)基本都是这种"薄封装 + 跨语言展开"的形态。例如 wordcount_textio.ts 展示了通过textio.readFromText("gs://dataflow-samples/shakespeare/kinglear.txt")读取 GCS 文本——底层正是经 Expansion Service 委托 Python SDK 的 TextIO 完成数据读取,示例中甚至直接演示了手动启动本地 Job Service 后以new PortableRunner("localhost:3333")对接的方式:
// python apache_beam/runners/portability/local_job_service_main.py --port 3333 await new PortableRunner("localhost:3333").run(async (root) => { const lines = await root.applyAsync( textio.readFromText("gs://dataflow-samples/shakespeare/kinglear.txt"), ); lines.apply(wordCount).map(console.log); });注意此处使用了root.applyAsync(...)——跨语言 Transform 需要异步展开,这正是 SDK 同时提供apply与applyAsync的原因。
五、API 设计:对 Beam 惯例的取舍
README-dev.md将 API 层面标记为仍在演进,而README.md(sdks/typescript/README.md)给出了更完整的取舍说明。总体原则是用 TypeScript 惯用法表达 Beam 概念,但不拘泥于传统 SDK 的形式,具体包括:
- 关系式基础(relational foundations):以带 schema 的数据为第一公民,使用 JavaScript 原生 Object 作为行类型(row type),弱化强 KV 型 Transform,转而用字段名或表达式定位数据;
- 淡化 Coder:Coder 降级为用于互操作的进阶特性;能从元素推断 schema 时就推断,否则使用基于 BSON 编码的兜底 Coder;
- PCollection 增加
map/flatMap方法,而不是只允许apply;apply同时接受函数((PCollection) => ...)与 PTransform 子类; - 取消独立的 Pipeline 对象:以
RootPValue 作为 Pipeline 构建起点,直接在 Runner 上调用run()(pvalue.ts 中的Root类即是入口); - PValue 可以是数组或对象:用
P(...)操作符包裹(如P([pc1, pc2, pc3]).apply(new Flatten())),避免引入 PCollectionTuple/PCollectionList; - 生成器式多输出:
flatMap与ParDo.process通过yield产出多个元素;需要多路输出时,使用Split原语把PCollection<{a?, b, ...}>拆成{a: PCollection, b: PCollection, ...}; - 可选 context 参数:
map/flatMap/ParDo.process可携带附加 context 对象,其中成员要么是常量,要么是DoFnParam之类的特殊参数,在运行时提供元素级信息(时间戳、窗口、旁路输入等); - 异步优先:由于 JS 生态天然异步且无法从异步回到同步,SDK 提供
PValue.applyAsync(run/runAsync同理),用户回调的全面异步化仍在 TBD。
正是这些取舍,使得 TypeScript SDK 在保留 Beam 语义的同时,具备了不同于 Java/Python SDK 的轻量手感。
六、测试、代码风格与文档
6.1 运行测试
npm testpretest钩子会先自动执行npm run build,随后mocha dist/test dist/test/docs运行已编译的测试用例。仓库中已有较成体系的测试覆盖(test 目录):primitives_test、combine_test、io_test、coders系列测试(js_coders_test、row_coder_test、standard_coders_test)、serialize_test、worker_test,以及文档示例测试 test/docs/programming_guide.ts。
6.2 代码风格
SDK 采用 Prettier 统一格式,全量格式化:
npx prettier --write .提交前可用npm run prettier-check校验、npm run lint做 ESLint 检查。
6.3 生成 API 文档
npm run docs会先构建,再由typedoc依据 typedoc.json 的配置产出文档。
七、当前状态与已知 TODO(以仓库为准)
README-dev.md明确声明该 SDK 仍在持续演进(work in progress)。截至文档记录(2022 年 1 月已具备构建并运行基础 Pipeline 的能力,包括外部 Transform 与可移植运行器),剩余的大项工作包括:
容器化(Containerization)
- 真正使用 worker threads 并行处理多个 bundle(是否收益显著尚不确定,当前通过 sibling workers 缓解)。
API
- 多处小特性或设计决策待定:考虑以双数组(2-arrays)替代
{key, value}对象表示 KV;强制map/flatMap的第二参数为 Object 以避开与Array.map的混淆,并考虑增加doFilter/doReduce;逐步摆脱类(classes)风格; - 高级特性(state、timers、SDF)尚未实现。
其他
- 相对/绝对导入策略(可能通过
jsconfig.json的 baseUrl 解决); - 更多更好的测试,包括对非法/不支持用法的测试;
- 像其他 SDK 一样设置 gRPC channel 选项(如
grpc.max_{send,receive}_message_length); - 减少
any的使用(可用unknown替代真正未知类型;若生成 proto 文件能被忽略,则重新启用noImplicitAny: true); - 引入 ESLint 并至少修复低垂果实。
对照当前源码可以发现,部分 TODO 已取得进展:package.json中已包含eslint及其lint脚本,tsconfig.json也已开启strictNullChecks。其余项(worker threads、state/timers/SDF 高级特性、grpc.max_*channel 选项等)与文档描述一致,仍未落地。
八、开发速查清单
| 场景 | 命令 |
|---|---|
| 安装依赖 | cd sdks/typescript && npm install |
| 构建 | npm run build |
| 本地运行 wordcount | node dist/src/apache_beam/examples/wordcount.js --runner=direct |
| Flink 运行 | node dist/src/apache_beam/examples/wordcount.js --runner=flink |
| Dataflow 运行 | node dist/src/apache_beam/examples/wordcount.js --runner=dataflow --project=${PROJECT_ID} --tempLocation=gs://${GCS_BUCKET}/wordcount-js/temp --region=${REGION} |
| 运行测试 | npm test |
| 格式检查/格式化 | npm run prettier-check/npx prettier --write . |
| Lint | npm run lint |
| 生成文档 | npm run docs |
结语
Apache Beam TypeScript SDK 是一条"以小博大"的移植路线:它不重复实现全部 Beam 语义,而是依靠可移植性框架——Runner API protobuf、跨语言 Transform 与 Expansion Service、可复用 Worker、Job Service 驱动的可移植运行器——用最少的原生代码获得完整的 Beam 能力。对开发者而言,这意味着一份 TypeScript Pipeline 既可以本地直跑,也可以无缝迁移到 Flink、Dataflow 等生产环境;对 SDK 开发者而言,它则是一份"如何用新语言快速实现 Beam"的活的参考教材。结合本仓库的源码(尤其是runners/、worker/、transforms/external.ts三处)阅读本文,可以更直观地把握这条执行链路的每个环节。
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam TypeScript SDK 开发者指南:从本地构建到多 Runner 运行的完整实践
Apache Beam TypeScript SDK 开发者指南:从本地构建到多 Runner 运行的完整实践 本指南以 Apache Beam 仓库中 sdk
大数据批处理流处理数据工程Apache Zeppelin Beam 解释器指南:架构设计、构建方法与源码级运行原理
Apache Zeppelin Beam 解释器指南:架构设计、构建方法与源码级运行原理 本文以 beam/README.md https://link.git
后端前端大数据数据分析Apache Beam Python SDK 开发实战指南:环境搭建、测试、构建与流水线运行
Apache Beam Python SDK 开发实战指南:环境搭建、测试、构建与流水线运行 导读 :本文是面向 Apache Beam Python SDK(
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考