☰
Apache Beam TypeScript SDK 开发指南:从源码构建、运行 Pipeline 到可移植运行器的实现原理
2026/10/10 8:32:48 网站建设 项目流程

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam18/beam
点击查看免费下载

导读

本文面向希望以 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 生态的一次原生尝试,它承载着两个截然不同的目标:

  1. 触达庞大的 JavaScript 开发者社区。现有数据处理框架对 JavaScript 开发者的支持相对不足,一个原生针对该语言的 SDK 可以填补这一空白。
  2. 充当 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 install

npm 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 build

package.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 cleantsc --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(可移植)运行器
directDirectRunner(direct_runner.ts)
universaluniversalRunner:通过 Python 启动本地 Job Service 的可移植运行器(universal.ts)
flinkflinkRunner:自动拉起 Flink Job Server(flink.ts)
dataflowdataflowRunner:对接 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)可归纳为:

  1. 探测流式需求:扫描所有 PCollection,若存在isBounded == UNBOUNDED,自动把streaming选项置为 true;
  2. 选择执行环境:
    • LOOPBACK模式:通过ExternalWorkerPool(external_worker_service.ts 中的 gRPC ExternalWorkerPool 服务)在本地进程内启动 Worker,适合本地验证;
    • 默认模式:将 SDK 环境替换为 Docker 环境,镜像为docker.io/apache/beam_typescript_sdk:<版本>(可用sdkContainerImage选项覆盖),并自动执行npm pack把当前代码打成 npm 包作为 artifact 注册,同时收集file:形式的本地依赖一并上传;
  3. 注册模块:把需要 Worker 端 import 的模块集合写入registeredNodeModules(来自serialization.getRegisteredModules());
  4. Prepare + Run:调用 Job Service 的Prepare方法提交PrepareJobRequest(Pipeline 与转成beam:option:*前缀的 pipeline options);若响应要求暂存 artifact,则通过 Artifact Staging Service 上传;最后调用Run方法拿到jobId;
  5. 返回句柄: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 的展开过程为:

  1. 构造ExpansionRequest,携带要展开的PTransform(含urn与 payload 配置)及其输入 PCollection;
  2. 调用指定地址的 Expansion Service(gRPC)的expand方法;
  3. 若展开结果的环境带有依赖(dependencies),则通过 Artifact Retrieval Service 拉取这些 artifact 转存为持久形式(resolveArtifacts);
  4. 将返回的 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 test

pretest钩子会先自动执行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
本地运行 wordcountnode 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 .
Lintnpm 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.

项目地址:https://gitcode.com/gh_mirrors/beam18/beam
点击查看免费下载

相关推荐

上一篇:Poetry终极问题排查指南:10个常见错误和快速解决方案
下一篇:Langchain-Chatchat 0.3.x版本前瞻:大模型智能体的终极进化指南

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询