Apache Beam YAML:用声明式配置快速构建批处理与流式管道
批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载Apache Beam YAML 是集成在 Apache Beam Python SDK 中的声明式管道编写方式你可以用纯 YAML 文件描述读取、转换、聚合、写出等完整数据处理流程无需编写 Python 代码、无需熟悉 SDK 编程模型。本文以仓库中 Beam YAML 模块的真实源码与示例为支撑带你掌握chain线性管道、非线性管道、流式窗口化管道、providers自定义变换的编写方法以及通过python -m apache_beam.yaml.main一键运行的完整实战流程。为什么用 YAML 描述 Beam 管道传统 Apache Beam 管道需要编写 Python 代码继承DoFn、定义PTransform、处理PCollection与 coder 等细节。Beam YAML API 改变了这一体验它带来几个核心价值零代码门槛管道的描述与编写完全在 YAML 文件中完成用任何文本编辑器即可维护不需要编程经验或 SDK 熟练度直观的声明式语法将读什么、做什么转换、写到哪里以数据表mapping/list的形式逐层声明结构一目了然比直接操作 Beam protos 更易理解Beam 管道底层是 protobuf 表示的PipelineProto手写与阅读都很困难YAML 语义更贴近人的思维是更适合人类编辑的中间表示适合作为工具间传递的中间表示由于可读性与易编写性YAML 非常适合作为管道编写 GUI、血缘分析lineage analysis等工具与 Beam 运行时之间的桥接格式。在实现上Beam YAML 解析器目前集成于 Apache Beam Python SDK位于 sdks/python/apache_beam/yaml 目录下包含main.py命令行入口、yaml_transform.py管道展开核心、yaml_provider.py变换提供方机制、yaml_mapping.py/yaml_combine.py/yaml_join.py等各内置变换的归一化实现以及完整的单元测试与可直接运行的示例 YAML。认识 pipeline 声明的基本骨架一个 Beam YAML 管道文件的最外层是一个pipeline键其内部可以包含以下组成部分键作用type管道编排方式取chain线性隐式连接或composite显式命名连接不写时按composite处理transforms变换列表是管道的主体每个变换条目至少包含type并可带name、input、output、config、windowingsource/sink仅用于chain型管道分别表示首尾的读取与写出变换options管道级 Beam PipelineOptions会透传给底层beam.Pipeline见 main.pyproviders自定义变换提供方声明用于引入外部Java / Python 包 / 远程服务的变换实现每个变换的通用字段为type变换类型名如ReadFromCsv、Filter、Sql、WriteToJson、WindowInto、Combine等name给变换起名用于监控、调试以及在同一个管道中区分多个同类型变换重复名称会触发去重处理见 yaml_transform.pyconfig该变换的参数映射input/output显式声明输入来源与输出去向非线性管道必需windowing为该变换的输入及应用本身附加窗口策略仅chain的transforms中的变换支持。从源码实现看expand_pipeline会先把pipeline_spec[pipeline]包装为根composite再交给YamlTransform逐层展开为真实的PTransform并接入beam.Pipeline见 yaml_transform.py。线性管道用chain隐式串联输入输出对于一条数据流从上到下依次处理的线性管道Beam YAML 允许你完全省略显式的input/output连接只要把pipeline.type设为chain各变换就会按列表顺序自动首尾相连。此时管道开头的读取变换记为source可选也可直接放在transforms第一项管道末尾的写出变换记为sink可选也可放在transforms最后一项source与sink在内部会被分别插入到transforms列表的首尾其输入被标记为显式空输入见 yaml_transform.pychain本质上会被重写为一个所有输入输出都隐式的composite变换每个后续变换自动以上一个变换的输出为输入见 yaml_transform.py。下面这个示例完整复刻了官方指南中的典型场景从 CSV 文件读取数据按条件过滤行执行一条 SQL 聚合查询再把结果写入 JSON 文件pipeline: type: chain source: type: ReadFromCsv config: path: /path/to/input*.csv transforms: - type: Filter config: language: python keep: col3 100 - type: Sql name: MySqlTransform config: query: select col1, count(*) as cnt from PCOLLECTION group by col1 sink: type: WriteToJson config: path: /path/to/output.json要点说明ReadFromCsv的path支持 glob 通配符如上例的input*.csv也支持本地路径与gs://等远端文件系统路径仓库示例 simple_filter.yaml 中即使用了gs://apache-beam-samples/beam-yaml-blog/products.csv这样的公开数据集Filter的language: python声明使用 Python 表达式求值keep给出保留条件等价于 Python SDK 的beam.Filter除python外language还可以指定其他受支持的表达式方言如generic、sqlSql变换由 Java 侧的 SQL 扩展服务提供详见下文 providersSQL 中PCOLLECTION是当前输入集合的固定别名WriteToJson的path指定输出目录或文件前缀。仓库中的 chain 实战示例仓库 examples 目录提供了大量可直接运行的 YAML 管道。以 simple_filter_and_combine.yaml 为例它演示了比官方文档示例更完整的组合读取商品 CSV → 过滤出Electronics类目 → 按product_name分组并统计销量与总营收 → 写回 CSVpipeline: transforms: - type: ReadFromCsv name: ReadInputFile config: path: gs://apache-beam-samples/beam-yaml-blog/products.csv - type: Filter name: FilterWithCategory input: ReadInputFile config: language: python keep: category Electronics - type: Combine name: CountNumberSold input: FilterWithCategory config: group_by: product_name combine: num_sold: value: product_name fn: count total_revenue: value: price fn: sum - type: WriteToCsv name: WriteOutputFile input: CountNumberSold config: path: output这个示例同时展示了非线性composite管道的写法transforms中每个变换都通过nameinput显式指定输入来源input值即上游变换的name。Combine的group_by指定分组字段combine子映射逐字段声明聚合函数count、sum等最终得到Row(product_nameHeadphones, num_sold2, total_revenue119.98)形式的聚合结果。非线性管道显式命名输入嵌套 chain并非所有数据处理都是线性的。当管道出现分支、合并或多个输入源时就必须使用非线性composite编排所有input/output都必须显式命名每个变换的input引用上游变换的name或用变换名.输出名形式引用多输出变换的特定输出见 yaml_transform.py。非线性管道同样支持把chain作为嵌套结构你可以在composite的transforms中放入一个type: chain的子管道让子链内部隐式连接、对外仅暴露input与output。这等价于把一个链包装成可复用的复合变换是组织复杂管道的重要手法。流式管道与窗口化Beam YAML 同时支持批处理batch与流式streaming管道。流式场景下Beam 的窗口与触发windowing triggering能力完整可用有两种声明方式使用标准WindowInto变换显式在transforms中插入一个WindowInto条目并在其config.windowing中声明窗口策略直接在变换上附加windowing给chain中某个变换挂上windowing键解析器会把该窗口应用到该变换的输入以及变换本身——从实现看preprocess_windowing会在该变换的输入前自动插入对应的WindowInto或对无输入的读取类变换把窗口推送到其根输出见 yaml_transform.py。以下示例来自官方指南演示一条完整的流式链从 Pub/Sub 主题读取消息应用带滑窗的分组变换再写回另一个 Pub/Sub 主题pipeline: type: chain transforms: - type: ReadFromPubSub config: topic: myPubSubTopic format: ... schema: ... - type: SomeGroupingTransform config: arg: ... windowing: type: sliding size: 60s period: 10s - type: WriteToPubSub config: topic: anotherPubSubTopic format: json这里windowing声明了一个滑动窗口size: 60s表示窗口长度 60 秒period: 10s表示每 10 秒滑动一次type还支持fixed、global、session等 Beam 标准窗口类型。ReadFromPubSub的format与schema用于声明消息编码与字段结构WriteToPubSub的format: json指定输出为 JSON。窗口化处理的内部机制源码 yaml_transform.py 中的preprocess_windowing会把挂在变换上的windowing拆解为一条前置的WindowInto变换必要时还会把复合变换内的窗口策略下推到其内部的读取、创建等根操作从而保证窗口语义与 Beam 原生WindowInto完全一致。运行管道python -m apache_beam.yaml.main编写好 YAML 文件后使用 Python SDK 的标准python -m命令即可运行python -m apache_beam.yaml.main --yaml_pipeline_file/path/to/pipeline.yaml [other pipeline options such as the runner]命令入口位于 main.py其执行流程为读取 YAML 文件或直接使用--yaml_pipeline传入的 YAML 字符串→ 展开 Jinja 模板变量 → 用SafeLineLoader解析 YAML → 可选地做 JSON Schema 校验 → 构造beam.Pipeline并展开执行。支持的参数详解参数默认值说明--yaml_pipeline_file别名--pipeline_spec_file无指定包含管道 YAML 描述的文件路径支持本地路径与远端文件系统如gs://--yaml_pipeline别名--pipeline_spec无直接以字符串形式传入 YAML 管道描述与yaml_pipeline_file二选一同时给出会报错--json_schema_validationgenericnone完全不校验generic校验管道整体结构但不校验单个变换per_transform进一步校验已知变换的config需要先运行generate_yaml_docs生成transforms.schema.yaml见 yaml_transform.py--jinja_variables{}以 JSON 字典形式传入变量供 YAML 中 Jinja 模板预处理使用--jinja_variable_flags[]逗号分隔的旗标名列表这些旗标会被提升为 Jinja 变量便于 Dataflow 模板等工具以扁平旗标形式传参除上述 Beam YAML 专属参数外其余参数如--runnerDirectRunner、--runnerDataflowRunner、--project、--region、--streaming等会原样透传给底层PipelineOptions。管道级 options 与 Jinja 模板在 YAML 文件内部的顶层pipeline.options中也可以声明管道选项源码中通过SafeLineLoader.strip_metadata剥离行号等元数据后并入PipelineOptions见 main.py。此外YAML 管道在解析前会先经过 Jinja2 模板预处理expand_jinja见 yaml_transform.py你可以在 YAML 中使用{{ variable }}占位符并用--jinja_variables{variable: value}在运行期注入从而让同一份模板在不同环境开发/生产间复用。providers引入内置与自定义变换的机制Beam YAML SDK 自带丰富的内置变换IO、Filter、MapToFields、Combine、Sql、Join、Partition、WindowInto 等同时通过providers概念支持自定义变换。从源码 yaml_provider.py 看Provider抽象负责把变换类型名与参数映射为具体的PTransform实现运行时会在create_ptransform时选择可用的 provider并通过affinity启发式尽量让相邻变换落在同一执行环境中以利于融合见 yaml_transform.py。内置 provider 类型providertype说明javaJar通过本地 Java Jar 启动跨语言扩展服务ExternalJavaProvider适用于 Java 侧的 IO / SQL 等变换mavenJar指定group_id、artifact_id、version从 Maven 中央仓库下载 Jar 后启动扩展服务beamJar指定gradle_target使用当前 Beam 版本构建的 Jar 启动扩展服务如 SQL 扩展服务sdks:java:extensions:sql:expansion-service:shadowJarremote连接远程已部署的扩展服务地址address参数适合生产环境复用预启动服务python通过全限定名fully qualified name引用本地的 PythonPTransform构造器无需额外安装pythonPackage指定packages从 PyPI 安装后通过全限定名引用其中的变换例如前文Sql变换之所以可用正是因为内置的SqlBackedProvider默认关联了beam:external:java:sql:v1这个 Java 侧 SQL 扩展服务的 URN并用beamJar方式提供实现见 yaml_provider.py。在 YAML 中声明自定义 provider在pipeline顶层或composite/chain内部的providers列表中声明即可例如pipeline: providers: - type: beamJar transforms: MyCustomTransform: beam:transforms:custom:my_custom:v1 config: gradle_target: sdks:java:extensions:my-custom:expansion-service:shadowJar transforms: - type: MyCustomTransform config: arg: value这样MyCustomTransform就会成为本管道内可用的变换类型。通过 providers你可以把跨语言Java变换、第三方 Python 包中的变换以及远程扩展服务无缝接入 YAML 管道这正是 Beam YAML 扩展能力的核心。管道展开的源码级流程理解底层机制有助于排查问题。YamlTransform在展开一个管道时会依次执行多阶段的预处理见 yaml_transform.pyensure_transforms_have_types校验每个变换都有typenormalize_mapping/normalize_combine将MapToFields、Combine等变换的声明归一化为内部标准形式preprocess_languages根据config.languagepython/generic/sql等把Filter、MapToFields等类型改写为Filter-python、MapToFields-generic等具体方言类型preprocess_source_sink把source/sink折叠进transformspreprocess_chain把chain重写为隐式连接的compositepreprocess_flattened_inputs为多输入自动插入Flattenpreprocess_windowing把windowing声明拆解为WindowInto变换ensure_errors_consumed校验所有error_handling输出都被消费。此外解析器采用自定义的SafeLineLoader见 yaml_transform.py会给每个映射和字符串附加__line__行号与__uuid__标识——行号用于在 Schema 校验或运行出错时给出精确的YAML 第几行定位uuid用于在展开过程中追踪变换之间的连接关系。这也是为什么 Beam YAML 报错信息通常能直接指出 YAML 文件的具体位置。仓库中的可运行示例与测试想深入实践可以直接查看并运行仓库内现成的资源示例 YAMLsdks/python/apache_beam/yaml/examples 目录下包含simple_filter.yaml、simple_filter_and_combine.yaml、regex_matches.yaml以及transforms/子目录中覆盖aggregation、mapping、join、window、io等分类的大量最小可运行示例文档化示例与测试main_test.py 演示了通过--yaml_pipeline_file运行 YAML 管道的测试写法readme_test.py、programming_guide_test.py 则确保仓库内 YAML 文档与示例始终可运行生成变换文档python -m apache_beam.yaml.generate_yaml_docs --schema_file...可以为已知变换生成transforms.schema.yaml开启per_transform级别的config校验。小结Beam YAML 把 Apache Beam 的强大能力封装进了可读、可版本化、可被工具消费的 YAML 描述中chain让线性管道近乎零配置显式命名让非线性管道清晰可控windowing声明让流式窗口开箱即用providers让自定义与跨语言变换无缝接入。对于希望快速搭建数据处理管道、或需要以中间表示驱动管道编排工具的场景Beam YAML 提供了一条从声明到运行python -m apache_beam.yaml.main --yaml_pipeline_file...的完整路径。赞分享批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载相关推荐Apache Beam YAML 管线编程指南用声明式配置构建批处理与流式数据管道Apache Beam YAML 管线编程指南用声明式配置构建批处理与流式数据管道 Apache Beam 的 YAML API 允许开发者仅用一个 YAML大数据批处理流处理数据工程Apache Beam YAML 实战指南用声明式配置构建首个无代码数据管道Apache Beam YAML 实战指南用声明式配置构建首个无代码数据管道 Apache Beam YAML 是 Apache Beam 推出的首个无代码大数据批处理流处理数据工程Apache Beam Spark Runner 全面指南用 Apache Spark 执行 Beam 批处理与流式管道Apache Beam Spark Runner 全面指南用 Apache Spark 执行 Beam 批处理与流式管道 导读 Apache Beam 通过统大数据批处理流处理数据工程上一篇FreeMove智能文件迁移工具解决C盘空间不足与程序移动难题下一篇FreeMove终极教程3分钟学会安全迁移C盘文件释放空间创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考