ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Apache Beam 多语言流水线实战:Python 与 Java 跨语言 Transform 互调指南

Apache Beam 多语言流水线实战:Python 与 Java 跨语言 Transform 互调指南 【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载本文基于 Apache Beam 仓库中的 examples/multi-language/README.md 与配套示例源码系统讲解 Beam 多语言流水线Multi-language Pipelines的完整实战方案如何在 Python 流水线中调用 Java 变换含JavaExternalTransformAPI如何在 Java 流水线中调用 Python 变换scikit-learn MNIST 分类与 DataFrame 词频统计并给出从扩展开启、环境搭建到 Dataflow/DirectRunner 上运行、结果校验的全套命令。读完本文你将掌握 Beam 跨语言展开服务Expansion Service的启动方式、URN 注册与调用机制并能直接复现仓库中的全部多语言示例。一、多语言流水线是什么Apache Beam 采用统一的编程模型Pipeline 抽象为 PCollection 上的 PTransform但不同语言 SDK 的能力并不完全对等——例如 Python 生态中强大的 scikit-learn 推理、Java 生态中成熟的Count.perElement()等内置变换。多语言流水线通过**展开服务Expansion Service**打破语言边界主流水线在本地 SDK 中构造一个ExternalTransform把变换的 URN统一资源名与参数通过 gRPC 发送给另一语言的展开服务由后者展开成完整的 Beam 子图Subgraph再交回主流水线继续执行。本仓库的 examples/multi-language 目录正是这一机制的最小可运行示例集合覆盖双向调用Python 主流水线调用 Java 变换python/addprefix.py、python/javacount.py、python/javadatagenerator.pyJava 主流水线调用 Python 变换SklearnMnistClassificationscikit-learn MNIST 手写数字分类与PythonDataframeWordCountPython DataFrame 词频统计位于 examples/java/src/main/java/org/apache/beam/examples/multilanguage。二、在 Python 流水线中使用 Java 变换三个 Python 示例均以 Java 侧实现为核心Python 侧只负责 I/O 与结果格式化是理解跨语言调用的最佳入门材料。2.1 addprefix向字符串添加前缀addprefix.py 读取文本文件先由 Java 变换JavaPrefix为每行加上java:前缀再由 Python 侧的beam.Map加上python:前缀最终每行输出形如python:java:原文。核心调用代码java_output ( input | JavaPrefix beam.ExternalTransform( beam:transform:org.apache.beam:javaprefix:v1, ImplicitSchemaPayloadBuilder({prefix: java:}), (localhost:%s % expansion_service_port)))这里beam:transform:org.apache.beam:javaprefix:v1就是 Java 变换注册的 URN第三个参数指向展开服务地址。Java 侧对应实现为 JavaPrefix.java它是一个标准PTransformPCollectionString, PCollectionString通过DoFn对每个元素执行prefix inputclass AddPrefixDoFn extends DoFnString, String { ProcessElement public void process(Element String input, OutputReceiverString o) { o.output(prefix input); } }2.2 javacount用 Java 的 Count.perElement() 统计词频javacount.py 先用 Python 的WordExtractingDoFn提取单词然后把单词 PCollection 交给 Java 变换JavaCount完成去重计数最后在 Python 侧格式化为word:count输出java_output ( words | JavaCount beam.ExternalTransform( beam:transform:org.apache.beam:javacount:v1, None, (localhost:%s % expansion_service_port)))Java 侧 JavaCount.java 的实现极其精简直接复用 Java SDK 内置变换Override public PCollectionKVString, Long expand(PCollectionString input) { return input.apply(JavaCount, Count.perElement()); }2.3 javadatageneratorJavaExternalTransform API 免注册调用与前两个示例不同javadatagenerator.py 演示了JavaExternalTransformAPI——通过类名直接引用 Java 变换无需预先在展开服务中注册 URNjava_transform JavaExternalTransform( org.apache.beam.examples.multilanguage.JavaDataGenerator, expansion_service(localhost:%s % expansion_service_port)).create(np.int32(100)).withDataConfig(data_config) data p | Generate java_transform其中DataConfig是一个具名元组用于传入 Java 侧withDataConfig的配置对象prefixstart、length20、suffixend。Java 侧 JavaDataGenerator.java 是无输入PBegin的数据生成变换用GenerateSequence.from(0).to(size)生成序列再经GenerateDataDoFn输出prefix *****... suffix形式的字符串星号个数由length扣除前后缀长度后决定。该示例清楚地展示了 Java 变换的构造器方法create 链式配置withDataConfig如何被跨语言暴露给 Python。2.4 Java 侧如何注册一个可跨语言调用的变换从源码结构看Java 侧一个可被 Python 调用的变换需要三个协作组件以JavaPrefix为例组件文件职责变换本体JavaPrefix.java实现PTransform的实际逻辑配置类JavaPrefixConfiguration.java承载外部传入的参数如prefixBuilderJavaPrefixBuilder.java实现ExternalTransformBuilderConfig, In, Out把配置转成变换实例RegistrarJavaPrefixRegistrar.java用AutoService(ExternalTransformRegistrar.class)注册 URN → Builder 映射Registrar 的核心是 URN 映射表AutoService(ExternalTransformRegistrar.class) public class JavaPrefixRegistrar implements ExternalTransformRegistrar { static final String URN beam:transform:org.apache.beam:javaprefix:v1; Override public MapString, ExternalTransformBuilder?, ?, ? knownBuilderInstances() { return ImmutableMap.of(URN, new JavaPrefixBuilder()); } }Python 侧beam.ExternalTransform传入的 URN 必须与 Registrar 中声明的 URN 完全一致JavaExternalTransform则绕开 URN 注册表直接按全限定类名展开。两个机制各有适用场景URN 方式适合长期稳定的公共变换类名方式适合快速原型。三、Python 侧运行实战从 JAR 到流水线3.1 启动 Java 展开服务先下载beam-examples-multi-languageJAR从 Apache Beam 2.36.0 起可在 Maven Central 按org.apache.beam组坐标检索获取然后启动展开服务java -jar beam-examples-multi-language-version.jar port --javaClassLookupAllowlistFile*参数说明port展开服务监听端口后续所有 Python 示例都通过--expansion_service_port PORT引用它--javaClassLookupAllowlistFile*放开JavaExternalTransform的类查找白名单javadatagenerator示例依赖此配置。3.2 准备 Python 运行环境按 Apache Beam Python 快速入门创建并激活虚拟环境然后安装 Apache Beam Python SDK按需加上[gcp]扩展以支持 Dataflow。若使用 DirectRunner 且跨语言环境需要隔离还需准备 Docker 环境并指定--environment_typeDOCKER。3.3 执行 Python 流水线在新的 Shell 中使用支持多语言流水线的 Runner 运行python/目录下的示例。三个示例文件头部均内置了完整的运行命令这里摘录如下。DirectRunner本地验证python addprefix.py --runner DirectRunner --environment_typeDOCKER \ --input INPUT FILE --output output --expansion_service_port PORTDataflowRunner云端运行python addprefix.py \ --runner DataflowRunner \ --temp_location $TEMP_LOCATION \ --project $GCP_PROJECT \ --region $GCP_REGION \ --job_name $JOB_NAME \ --num_workers $NUM_WORKERS \ --input gs://dataflow-samples/shakespeare/kinglear.txt \ --output gs://$GCS_BUCKET/javaprefix/output \ --expansion_service_port PORTjavacount.py与javadatagenerator.py的命令结构一致javadatagenerator无需--input因为 Java 变换自主生成数据。注意所有示例都要求--expansion_service_port指向第 3.1 步启动的展开服务。四、在 Java 流水线中使用 Python 变换反向场景同样被仓库覆盖Java 主流水线通过 Python 展开服务调用 Python 变换两个示例都位于 examples/java/src/main/java/org/apache/beam/examples/multilanguage。4.1 Sklearn MNIST 分类SklearnMnistClassification对 MNIST 手写数字图片进行图像分类读取 CSV含标签与像素喂给 Python 侧基于 scikit-learn 训练的分类模型输出每行由逗号分隔的真实标签,预测标签。准备阶段Setup准备输入 CSV 文件包含标签和像素上传至 GCS 存储桶准备 scikit-learn 模型的 pickle 文件基于 MNIST 数据训练上传至 GCS按所选 Runner 完成环境准备——下文以 Dataflow Runner 为例其他可移植 Runner 需按其指南修改参数。方式 A在已发布版本上运行Beam 2.43.0 及之后先用 Beam 官方 Maven 原型生成示例工程export BEAM_VERSIONBeam version mvn archetype:generate \ -DarchetypeGroupIdorg.apache.beam \ -DarchetypeArtifactIdbeam-sdks-java-maven-archetypes-examples \ -DarchetypeVersion$BEAM_VERSION \ -DgroupIdorg.example \ -DartifactIdmulti-language-beam \ -Dversion0.1 \ -Dpackageorg.apache.beam.examples \ -DinteractiveModefalse然后运行流水线export GCP_PROJECTGCP project export GCP_BUCKETGCP bucket export GCP_REGIONGCP region mvn compile exec:java -Dexec.mainClassorg.apache.beam.examples.multilanguage.SklearnMnistClassification \ -Dexec.args--runnerDataflowRunner --project$GCP_PROJECT \ --region$GCP_REGION \ --gcpTempLocationgs://$GCP_BUCKET/multi-language-beam/tmp \ --outputgs://$GCP_BUCKET/multi-language-beam/output \ -Pdataflow-runner校验输出每行第一个字段是数字真实标签第二个字段是预测标签gsutil cat gs://$GCP_BUCKET/multi-language-beam/output*方式 B在仓库 HEAD 上运行Beam 2.41.0 / 2.42.0 场景此方式需要自行构建 Python 与 Java 的 SDK 容器。步骤依次为按 Beam Python 快速入门创建并激活虚拟环境安装带 GCP 支持的 Beam 包与sklearnpip install apache-beam[gcp] pip install sklearn启动 Python 展开服务注意这里用--fully_qualified_name_glob *放开类名匹配对应 JAR 侧的白名单参数python -m apache_beam.runners.portability.expansion_service_main -p PORT --fully_qualified_name_glob *确认本机已安装 Docker然后在另一个 Shell 构建并推送 Python 与 Java 的 SDK 容器export DOCKER_ROOTDocker root ./gradlew :sdks:python:container:py38:docker -Pdocker-repository-root$DOCKER_ROOT -Pdocker-taglatest docker push $DOCKER_ROOT/beam_python3.8_sdk:latest ./gradlew :sdks:java:container:java11:docker -Pdocker-repository-root$DOCKER_ROOT -Pdocker-taglatest -Pjava11Home$JAVA_HOME docker push $DOCKER_ROOT/beam_java11_sdk:latest用 Gradle 运行流水线同时覆盖 Java 与 Python 两个 SDK Harness 容器并指向本地展开服务export GCP_PROJECTGCP project export GCP_BUCKETGCP bucket export GCP_REGIONGCP region export EXPANSION_SERVICE_PORTPORT # 清理可能存在的旧输出 gsutil rm gs://$GCP_BUCKET/multi-language-beam/output* ./gradlew :examples:multi-language:sklearnMinstClassification --args \ --runnerDataflowRunner \ --project$GCP_PROJECT \ --gcpTempLocationgs://$GCP_BUCKET/multi-language-beam/tmp \ --outputgs://$GCP_BUCKET/multi-language-beam/output \ --sdkContainerImage$DOCKER_ROOT/beam_java11_sdk:latest \ --sdkHarnessContainerImageOverrides.*python.*,$DOCKER_ROOT/beam_python3.8_sdk:latest \ --expansionServicelocalhost:$EXPANSION_SERVICE_PORT \ --region${GCP_REGION}参数要点--sdkContainerImage指定 Java SDK Harness 容器镜像--sdkHarnessContainerImageOverrides.*python.*,...用正则把 Python SDK Harness 容器替换为自建镜像多语言场景的必备参数--expansionServicePython 展开服务地址与第 3 步的PORT对应。输出校验命令与方式 A 相同gsutil cat gs://$GCP_BUCKET/multi-language-beam/output*4.2 Python Dataframe WordcountPythonDataframeWordCount在 Java 流水线中通过PythonExternalTransform调用 Python 的DataframeTransform完成词频统计。其文件头部给出了 Dataflow 上的运行示例./gradlew :examples:multi-language:pythonDataframeWordCount --args \ --runnerDataflowRunner \ --outputgs://{$OUTPUT_BUCKET}/count \ --sdkHarnessContainerImageOverrides.*python.*,gcr.io/apache-beam-testing/beam-sdk/beam_python{$PYTHON_VERSION}_sdk:latest源码可见该示例完整演示了 Java 侧调用 Python 变换的标准姿势通过 PythonExternalTransform 指定 Python 侧变换并把 Java 的Row数据与 Python DataFrame 通过 Beam Schema 打通。Java 侧还通过ExtractWordsFnDoFnString, Row把文本切词后转为带 Schema 的Row体现了多语言流水线中Schema 即契约的数据交换方式。五、两种方向的运行差异速查对比维度Python 调 JavaJava 调 Python展开服务java -jar beam-examples-multi-language-version.jar portpython -m apache_beam.runners.portability.expansion_service_main -p PORT客户端 APIbeam.ExternalTransform(URN, payload, service)/JavaExternalTransform(类名, ...)PythonExternalTransformJava SDKextensions.python包参数放开--javaClassLookupAllowlistFile*--fully_qualified_name_glob *典型场景复用 Java 内置/成熟变换如Count.perElement()复用 Python 机器学习/DataFrame 生态如 sklearn、pandas容器要求DirectRunner 下可加--environment_typeDOCKER需构建并覆盖 Python/Java SDK Harness 容器镜像六、结语与进一步阅读多语言流水线让团队不必为某种语言独有能力重写业务逻辑本文四个示例分别覆盖了 Python→Java 的 URN 注册调用、免注册类名调用以及 Java→Python 的模型推理与 DataFrame 处理。动手实践时建议先在本地按第 3 节启动展开服务并用 DirectRunner 跑通addprefix.py再逐步迁移到 Dataflow 上验证容器覆盖与--expansionService参数。进一步深入可以阅读Python 侧三个示例的完整命令注释addprefix.py、javacount.py、javadatagenerator.pyJava 侧变换注册全链路JavaPrefixRegistrar.java、JavaCountRegistrar.javaJava 调用 Python 的完整实现SklearnMnistClassification.java、PythonDataframeWordCount.java。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐Apache Beam 多语言流水线实战examples/multi-language 中 Python 与 Java 转换互相调用的完整指南Apache Beam 多语言流水线实战examples/multi language 中 Python 与 Java 转换互相调用的完整指南 Apache大数据批处理流处理数据工程Apache Beam跨语言转换开发Java与Python互操作Apache Beam跨语言转换开发Java与Python互操作 在企业数据处理场景中技术团队常面临多语言技术栈整合难题Java拥有成熟的大数据处理库P批处理流处理大数据Apache Fury跨语言序列化终极指南Java、Python、Go多语言数据互通实战Apache Fury跨语言序列化终极指南Java、Python、Go多语言数据互通实战 Apache Fury是一款高性能的跨语言序列化框架专为现代微服务基础架构开发工具上一篇FanControl中文设置终极指南5分钟让Windows风扇控制彻底汉化下一篇Bebas Neue字体完整指南免费开源标题字体的终极解决方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表