
大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载导读本文基于 Apache Beam 仓库中的官方示例 include/README.md 及配套文件深入讲解如何利用 Jinja 的% include指令将一个完整的 WordCount 流水线拆分为主流水线 子模块的可复用结构并通过--jinja_variables在运行时注入参数最终在 Dataflow或本地/其他 Runner上运行。读完本文你将掌握 Beam YAML 与 Jinja 模板结合的核心用法子模块拆分、变量插值、命令行 JSON 传参以及如何基于仓库源码理解 Jinja 渲染与参数解析的底层机制。一、示例背景为什么需要 Jinja% includeApache Beam 的 YAML 方言Beam YAML允许用纯声明式 YAML 描述整条流水线pipeline→type: chain→transforms无需编写 Python 代码。但当流水线变长、多个项目需要复用时单一 YAML 文件会变得难以维护。仓库在 sdks/python/apache_beam/yaml/examples/transforms/jinja/ 目录下提供了三种 Jinja 组织流水线的方式include/用% include指令将每个 transform 拆成独立的子模块文件主流水线通过相对路径引入import/用% import与宏macro复用 transform 片段见 import/README.mdinheritance/通过基流水线 派生覆盖的方式复用见 inheritance/README.md。本文聚焦include/模式一个主流水线 每个 transform 一个子模块。这样每个 transform 可以独立维护、单独测试、跨流水线复用同时主文件保持骨架化、可读性极高。本示例对应的目录结构为sdks/python/apache_beam/yaml/examples/transforms/jinja/include/ ├── README.md ├── wordCountInclude.yaml # 主流水线 └── submodules/ ├── readFromTextTransform.yaml # 读取文本 ├── mapToFieldsSplitConfig.yaml # 切词映射 ├── explodeTransform.yaml # 展开词数组 ├── combineTransform.yaml # 按词分组计数 ├── mapToFieldsCountConfig.yaml # 格式化输出 └── writeToTextTransform.yaml # 写出结果二、主流水线wordCountInclude.yaml 的骨架结构主流水线文件 wordCountInclude.yaml 定义了流水线的整体形态读取文本 → 切词 → 展开 → 计数 → 格式化 → 写出。注意它的 transforms 列表本身只写了两个内联 transformSplit words、Format output其余全部通过% include引入pipeline: type: chain transforms: # Read in text file {% include apache_beam/yaml/examples/transforms/jinja/include/submodules/readFromTextTransform.yaml %} # Split words and count occurrences - name: Split words type: MapToFields config: {% include apache_beam/yaml/examples/transforms/jinja/include/submodules/mapToFieldsSplitConfig.yaml %} # Explode into individual words {% include apache_beam/yaml/examples/transforms/jinja/include/submodules/explodeTransform.yaml %} # Group by word {% include apache_beam/yaml/examples/transforms/jinja/include/submodules/combineTransform.yaml %} # Format output to a single string consisting of word - count - name: Format output type: MapToFields config: {% include apache_beam/yaml/examples/transforms/jinja/include/submodules/mapToFieldsCountConfig.yaml %} # Write to text file on GCS, locally, etc {% include apache_beam/yaml/examples/transforms/jinja/include/submodules/writeToTextTransform.yaml %}2.1 整段 include把完整 transform 嵌进列表readFromTextTransform.yaml、explodeTransform.yaml、combineTransform.yaml、writeToTextTransform.yaml这四个子模块各自包含完整的 YAML transform 条目- name: ...开头。例如 readFromTextTransform.yaml- name: Read from GCS type: ReadFromText config: path: {{readFromTextTransform.path}}include 的结果会直接插入到 transforms 列表的对应位置等价于在主文件里写了- name: Read from GCS type: ReadFromText config: path: 运行时注入的路径2.2 局部 include只嵌入 config 片段mapToFieldsSplitConfig.yaml与mapToFieldsCountConfig.yaml则是只 include config 内容的子模块被嵌套在Split words、Format output这两个内联 transform 的config:下。例如 mapToFieldsSplitConfig.yamllanguage: {{mapToFieldsSplitConfig.language}} fields: word: callable: |- # TODO(#35936): Including another file here works fine, but if # the file has a license header or other irrevalent comments, it # will break the pipeline. Need to investigate more on Jinja # filtering in the expand_jinja method or some other way. import re def my_mapping(row): return re.findall(r[A-Za-z\], row.line.lower()) value: {{mapToFieldsSplitConfig.fields.value}}这里word字段通过 Python callable 对每行文本做正则切词re.findall(r[A-Za-z\], row.line.lower())提取全部字母含撇号构成的词并转为小写value字段被赋为常量1为后续计数做准备。文件中的注释还记录了一个已知注意点TODO(#35936)被 include 的文件若含许可证头或无关注释可能破坏流水线解析这正是仓库中对 Jinja 过滤/渲染细节仍需进一步调研的真实写照。类似地mapToFieldsCountConfig.yaml 只含 config 片段language: {{mapToFieldsCountConfig.language}} fields: output: {{mapToFieldsCountConfig.fields.output}}output字段的表达式在运行时由--jinja_variables注入本例为word - str(value)即把词与计数拼接成word - count字符串。2.3 剩余的两个子模块explodeTransform.yaml 将Split words产生的词数组逐个展开为独立元素- name: Explode word arrays type: Explode config: fields: - {{explodeTransform.fields}}combineTransform.yaml 按词分组并对value求和完成词频统计- name: Count words type: Combine config: group_by: - {{combineTransform.group_by}} combine: value: {{combineTransform.combine.value}}writeToTextTransform.yaml 将结果写出- name: Write to GCS type: WriteToText config: path: {{writeToTextTransform.path}}三、运行时参数注入--jinja_variables 详解被 include 的子模块通过{{变量名}}占位符引用参数这些变量由命令行--jinja_variables以 JSON 字典形式在运行时注入从而做到模板与参数分离。3.1 环境准备来自原文档原 README 给出的通用环境变量设置如下export PIPELINE_FILEapache_beam/yaml/examples/transforms/jinja/include/wordCountInclude.yaml export KINGLEARgs://dataflow-samples/shakespeare/kinglear.txt export TEMP_LOCATIONgs://MY-BUCKET/wordCounts/ export PROJECTMY-PROJECT export REGIONMY-REGION cd PATH_TO_BEAM_REPO/beam/sdks/python说明PIPELINE_FILE指向示例主流水线若从仓库根目录运行也可直接写sdks/python/apache_beam/yaml/examples/...。KINGLEAR是 Google Cloud 公开的莎士比亚《李尔王》文本文件Dataflow 示例数据集需能访问 GCS本地运行时可换成任意本地文本路径。TEMP_LOCATION为输出目录GCS 或本地需替换MY-BUCKET为实际 bucket。PROJECT/REGION对应 Dataflow 的 GCP 项目与区域仅在 Dataflow Runner 上需要本地/Direct Runner 可省略。从sdks/python目录运行是为了让 Python 包apache_beam处于可导入路径中。3.2 多行运行示例原文档原文python -m apache_beam.yaml.main \ --project${PROJECT} \ --region${REGION} \ --yaml_pipeline_file${PIPELINE_FILE} \ --jinja_variables{ readFromTextTransform: {path: ${KINGLEAR}}, mapToFieldsSplitConfig: { language: python, fields: { value: 1 } }, explodeTransform: {fields: word}, combineTransform: { group_by: word, combine: {value: sum} }, mapToFieldsCountConfig: { language: python, fields: {output: word \ - \ str(value)} }, writeToTextTransform: {path: ${TEMP_LOCATION}} }注意 JSON 内嵌 shell 变量的写法${KINGLEAR}表示闭合单引号、展开 shell 变量、再重新打开单引号最终拼出合法的 JSON 字符串。若本地运行可将${KINGLEAR}直接替换为/path/to/local/file.txt。3.3 单行运行示例原文档原文python -m apache_beam.yaml.main --project${PROJECT} --region${REGION} \ --yaml_pipeline_file${PIPELINE_FILE} --jinja_variables{readFromTextTransform: {path: ${KINGLEAR}}, mapToFieldsSplitConfig: {language: python, fields:{value:1}}, explodeTransform:{fields:word}, combineTransform:{group_by:word, combine:{value:sum}}, mapToFieldsCountConfig:{language: python, fields:{output:word \ - \ str(value)}}, writeToTextTransform:{path:${TEMP_LOCATION}}}单行与多行完全等价仅在 shell 书写习惯上不同适合放入脚本或 CI 配置。3.4 参数与占位符的映射关系--jinja_variables中的每个顶层键对应主流水线中一个子模块的命名空间键内字段对应子模块里的{{命名空间.字段}}占位符jinja_variables 顶层键注入的子模块占位符注入值示例readFromTextTransform.pathreadFromTextTransform.yaml 中{{readFromTextTransform.path}}gs://dataflow-samples/shakespeare/kinglear.txtmapToFieldsSplitConfig.language/.fields.valuemapToFieldsSplitConfig.yaml 中{{mapToFieldsSplitConfig.language}}、{{mapToFieldsSplitConfig.fields.value}}python/1explodeTransform.fieldsexplodeTransform.yaml 中{{explodeTransform.fields}}wordcombineTransform.group_by/.combine.valuecombineTransform.yaml 中{{combineTransform.group_by}}、{{combineTransform.combine.value}}word/summapToFieldsCountConfig.language/.fields.outputmapToFieldsCountConfig.yaml 中{{mapToFieldsCountConfig.language}}、{{mapToFieldsCountConfig.fields.output}}python/word - str(value)writeToTextTransform.pathwriteToTextTransform.yaml 中{{writeToTextTransform.path}}gs://MY-BUCKET/wordCounts/主流水线文件中内联的部分Split words、Format output的type、name是固定的只有通过占位符引用的值需要注入。四、Jinja 渲染与参数解析的底层实现4.1 expand_jinja渲染入口从源码看Beam YAML 的 Jinja 渲染入口位于 yaml_transform.py 的 expand_jinjadef expand_jinja( jinja_template: str, jinja_variables: Mapping[str, Any], search_paths: Iterable[str] ()) - str: beam_root_dir os.path.dirname(os.path.dirname(os.path.abspath(beam.__file__))) all_search_paths list(search_paths) if beam_root_dir not in all_search_paths: all_search_paths.append(beam_root_dir) if . not in all_search_paths: all_search_paths.append(.) return (jinja2.Environment( undefinedjinja2.StrictUndefined, loader_BeamFileIOLoader(all_search_paths)) .from_string(strip_leading_comments(jinja_template)) .render(datetimedatetime, **jinja_variables))这段实现揭示了几个关键行为搜索路径自动扩展除命令行传入的 search_paths来自yaml_pipeline_file所在目录见 main.py 的 _build_pipeline_yaml_from_argv外还会自动追加 Beam 根目录与当前目录。这正是示例中 include 路径写作apache_beam/yaml/examples/...的根目录全路径仍能解析成功的原因——从 Beam 包根目录开始定位。StrictUndefined使用jinja2.StrictUndefined任何未注入的变量都会立即抛错防止变量拼写错误导致静默渲染为空的隐患是模板化流水线的一种保护机制。strip_leading_comments渲染前会先剥离模板顶部的注释许可证头等这也呼应了 mapToFieldsSplitConfig.yaml 注释里 TODO(#35936) 提到的问题——被 include 文件内部的许可证注释仍需进一步处理。内置 datetime渲染环境还注入了datetime便于在模板中生成时间戳等动态值。4.2 main.py 的参数管线命令入口 main.py 的处理顺序是_preparse_jinja_flagsmain.py#L41-L87支持把任意--flag提升为 Jinja 变量通过--jinja_variable_flags指定白名单便于 Dataflow 模板等只能传扁平参数的工具使用若与已有 pipeline option 冲突则跳过。_parse_argumentsmain.py#L90-L135解析--yaml_pipeline/--yaml_pipeline_file互为别名--pipeline_spec/--pipeline_spec_file、--json_schema_validation、--jinja_variables、--tests等参数。注意--jinja_variables用json.loads解析所以必须是合法 JSON。_build_pipeline_yaml_from_argvmain.py#L244-L255读取模板 → 追加搜索路径 → 调用expand_jinja渲染出最终流水线 YAML。build_pipeline_components_from_yamlmain.py#L283-L303用SafeLineLoader加载渲染后的 YAML再交给expand_pipelineyaml_transform.py#L1515逐 transform 展开执行。渲染产物还会被放入display_datayaml、yaml_jinja_template、yaml_jinja_variables便于在 Runner UI / 日志中追溯模板长什么样、注入了哪些变量。4.3 测试验证仓库的测试 examples_test.py 展示了同样的渲染思路从测试数据中读取word_count_jinja_parameter_data()作为jinja_variables直接对模板调用template.render(jinja_variables)断言渲染结果。这证明Jinja 模板可以脱离 Beam Runner 独立渲染和单测主流水线与参数完全解耦。五、预期输出与运行说明主流水线注释给出了《李尔王》文本的预期输出示例wordCountInclude.yaml#L56-L66# Expected: # Row(outputking - 311) # Row(outputlear - 253) # Row(outputdramatis - 1) # Row(outputpersonae - 1) # Row(outputof - 483) # Row(outputbritain - 2) # Row(outputfrance - 32) # Row(outputduke - 26) # Row(outputburgundy - 20) # Row(outputcornwall - 75)每条输出都是word - count格式的字符串即mapToFieldsCountConfig.fields.output表达式word - str(value)的作用结果。运行前提与限制读取gs://dataflow-samples/shakespeare/kinglear.txt需要配置 Google Cloud 应用默认凭据ADC本地测试时建议把readFromTextTransform.path指向本地文件。使用 Dataflow Runner 时需要--project、--region与有效的TEMP_LOCATION使用本地 Direct Runner 时可去掉这些参数输出路径写本地目录。--jinja_variables必须是合法 JSON 字典且键需与模板占位符命名空间精确匹配否则 StrictUndefined 会直接报错。六、小结include 模式的最佳实践从本示例可以提炼出 Beam YAML Jinja% include的几条实用经验按 transform 粒度拆分每个子模块只含一个 transform整段 include或只含 config局部 include主文件保持骨架结构职责清晰。参数全部走运行时注入路径、语言、表达式等可变项统一用{{命名空间.字段}}占位符通过--jinja_variables注入实现同一模板、多套参数的复用。充分利用自动搜索路径include 路径既可用相对路径也可像示例一样写从 Beam 根目录出发的绝对包路径expand_jinja 会附加 Beam 根目录与当前目录作为搜索路径。善用 StrictUndefined 校验漏注入的变量会显式报错而不是静默消失建议在 CI 中用--testsmain.py#L112-L117配合测试套件对渲染后的流水线做回归验证。如果想进一步对比 Jinja 复用流水线的其他组织方式可以继续阅读同目录下的 import/README.md宏复用与 inheritance/README.md继承覆盖以及 examples/README.md 了解整个 YAML 示例集的运行方式。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam Python YAML SDK 的 Jinja2 % import 宏用宏文件复用流水线变换与配置Apache Beam Python YAML SDK 的 Jinja2 % import 宏用宏文件复用流水线变换与配置 导读 本文围绕 Apache Be大数据批处理流处理数据工程Apache Beam YAML 管线示例集解析从零样本到聚合、IO、Jinja 与 ML 管线的完整实战指南Apache Beam YAML 管线示例集解析从零样本到聚合、IO、Jinja 与 ML 管线的完整实战指南 Beam YAML 允许你用声明式 YAML大数据批处理流处理数据工程使用 Apache Beam YAML SDK 构建信用卡欺诈检测 MLOps 流水线特征工程与模型评估实战使用 Apache Beam YAML SDK 构建信用卡欺诈检测 MLOps 流水线特征工程与模型评估实战 Apache Beam YAML SDK 允许开大数据批处理流处理数据工程上一篇Dialog高级功能整合如何将对话框工具集成到大型项目中下一篇Flink SQL窗口函数性能飙升3个拆分预聚合黑科技创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考