ARTICLE · INTELLIGENCE

战地情报 · 详情页

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

Apache Beam TypeScript SDK 开发指南:从源码构建、运行 Pipeline 到可移植运行器的实现原理

Apache Beam TypeScript SDK 开发指南:从源码构建、运行 Pipeline 到可移植运行器的实现原理 【免费下载链接】beamApache 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、跨语言 Transformcross-language transforms经 Expansion Service 展开、可移植运行器经 Job Service 提交作业等底层机制。一、TypeScript SDK 的双重使命Apache Beam 是一个面向批处理与流处理的统一编程模型官方 SDK 覆盖 Java、Python 和 Go。而sdks/typescript目录下的 TypeScript SDKsdks/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 大量使用跨语言 TransformTypeScript 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等与开发依赖typescript4.7、mocha、eslint、prettier、typedoc等。2.2 构建TypeScript 编译为 JavaScriptnpm 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 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 生成覆盖率并上传 codecovnpm run worker直接以node运行外部 Worker 服务external_worker_service.js三、运行一个真实 Pipelinewordcountsdks/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 --runnerdirect--runnerdirect对应 sdks/typescript/src/apache_beam/runners/direct_runner.ts 中的directRunner工厂函数。它不依赖任何外部服务适合快速验证 Pipeline 逻辑。3.2 在 Flink 上运行基础设施自动下载node dist/src/apache_beam/examples/wordcount.js --runnerflink--runnerflink会走 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 \ --runnerdataflow \ --project${PROJECT_ID} \ --tempLocationgs://${GCS_BUCKET}/wordcount-js/temp --region${REGION}--runnerdataflow走 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.PCollectionstring): beam.PCollectionany { 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 到可移植 Runner4.1 运行器的创建与选择所有运行器都通过 sdks/typescript/src/apache_beam/runners/runner.ts 中的createRunner(options)统一创建支持以下取值--runner值实现default缺省defaultRunner能在 Direct 上跑的优先用 Direct否则回退到 universal可移植运行器directDirectRunnerdirect_runner.tsuniversaluniversalRunner通过 Python 启动本地 Job Service 的可移植运行器universal.tsflinkflinkRunner自动拉起 Flink Job Serverflink.tsdataflowdataflowRunner对接 Google Cloud Dataflowdataflow.tsRunner抽象类提供了两种执行入口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 自动交给universalRunnerenvironmentType: 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 触发一个元素而不是每个 workerDirectGbkOperator在单 bundle 内完成 GroupByKey对 key 按 windowkey 分组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模式通过ExternalWorkerPoolexternal_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方法提交PrepareJobRequestPipeline 与转成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 ServiceULRUniversal Local RunnerflinkRunner用JavaJarService启动 Flink Job Server jardataflowRunner用PythonService.forModule(apache_beam.runners.dataflow.dataflow_job_service, ...)启动 Dataflow Job Service。4.5 跨语言 TransformIO 的主力实现方式由于 TypeScript SDK 将 IO 大量委托给跨语言 Transform理解 Expansion Service 机制是掌握该 SDK 的关键。以 sdks/typescript/src/apache_beam/transforms/external.ts 中的rawExternalTransform为例跨语言 Transform 的展开过程为构造ExpansionRequest携带要展开的PTransform含urn与 payload 配置及其输入 PCollection调用指定地址的 Expansion ServicegRPC的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.mdsdks/typescript/README.md给出了更完整的取舍说明。总体原则是用 TypeScript 惯用法表达 Beam 概念但不拘泥于传统 SDK 的形式具体包括关系式基础relational foundations以带 schema 的数据为第一公民使用 JavaScript 原生 Object 作为行类型row type弱化强 KV 型 Transform转而用字段名或表达式定位数据淡化 CoderCoder 降级为用于互操作的进阶特性能从元素推断 schema 时就推断否则使用基于 BSON 编码的兜底 CoderPCollection 增加map/flatMap方法而不是只允许applyapply同时接受函数(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.applyAsyncrun/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本地运行 wordcountnode dist/src/apache_beam/examples/wordcount.js --runnerdirectFlink 运行node dist/src/apache_beam/examples/wordcount.js --runnerflinkDataflow 运行node dist/src/apache_beam/examples/wordcount.js --runnerdataflow --project${PROJECT_ID} --tempLocationgs://${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三处阅读本文可以更直观地把握这条执行链路的每个环节。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐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大数据批处理流处理数据工程上一篇Poetry终极问题排查指南10个常见错误和快速解决方案下一篇Langchain-Chatchat 0.3.x版本前瞻大模型智能体的终极进化指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

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