批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载本篇技术指南聚焦 Apache Beam 官方文档中Pipeline option patterns的核心模式——在管道pipeline作业完成后通过ValueProvider接口追溯访问并记录运行时参数。文章以website/www/site/content/en/documentation/patterns/pipeline-options.md为主体骨架结合仓库内 Java 与 Python SDK 的完整示例代码Snippets.java、snippets.py以及ValueProvider的底层实现ValueProvider.java、value_provider.py帮助你彻底理解为什么ValueProvider参数只能在 Beam DAG 内部被读取、如何通过管道分支 占位元素 DoFn这一官方推荐的解决方案实现参数追溯记录以及该模式向外部数据库落盘等场景的扩展思路。阅读完本文你将能直接在真实项目中复制、改造这套 Java / Python 双语言方案。一、问题背景为什么要在作业完成后记录运行时参数Apache Beam 的 Pipeline Options 用于配置管道的方方面面例如选择哪个 runner 执行管道、runner 特有的配置、项目 ID、文件存储位置等。官方编程指南指出可以通过命令行以--optionvalue的格式统一解析并填充PipelineOptions详见 配置管道选项。在此基础上Beam 提供了ValueProvider接口允许把某个选项的值延迟到运行时runner 实际执行作业时才注入而不是在管道构建graph construction阶段就固化。这种模板化参数能力在 Dataflow 模板等场景中极为常见同一份管道代码通过不同的运行时参数驱动不同的作业。然而ValueProvider有一个天然限制你只能在 Beam DAG由 PTransform 构成的执行图内部读取它的值。也就是说管道构建完成后、作业运行完毕在 DAG 之外的普通代码上下文里ValueProvider.get()是无法访问到真实值的Java 侧会抛出IllegalStateException: Value only available at runtime, but accessed from a non-runtime contextPython 侧对应RuntimeValueProviderError。这正是 patterns/pipeline-options.md 想要解决的问题你可以用ValueProvider把运行时参数传入管道但只能在 Beam DAG 内部记录logging这些参数。官方给出的解决方案是为管道添加一个分支branch——用一个DoFn处理一个占位placeholder值在DoFn内部读取并记录所有ValueProvider运行时参数。这样日志记录动作本身就成了 DAG 的一部分从而绕开了DAG 外无法取值的限制。二、模式拆解分支 占位元素 DoFn整个模式由三个要素构成占位输入使用Create.of(1)Java或beam.Create([None])Python制造一个极小的、只含单个元素的PCollection作为触发日志逻辑的引信。该分支的代价可忽略不计不参与真实业务计算。日志 DoFn在该分支上应用ParDo/beam.ParDoDoFn的process方法内通过PipelineOptions或直接持有ValueProvider引用调用get()取出运行时值并LOG.info输出。主管道不受影响真正的业务逻辑示例中仅用Sum.integersGlobally()求和占位照常挂在主分支上两个分支由同一个Pipeline.run()驱动。其精妙之处在于DoFn的process方法是在runner 的 worker 上、作业运行期间执行的此时RuntimeValueProvider的运行时上下文已就绪取值合法且必然成功。三、Java 实现完整代码与逐行讲解以下完整代码来自仓库示例 examples/java/src/main/java/org/apache/beam/examples/snippets/Snippets.java对应文档中的AccessingValueProviderInfoAfterRunSnip1片段。3.1 定义带 ValueProvider 的 PipelineOptions/** Sample of PipelineOptions with a ValueProvider option argument. */ public interface MyOptions extends PipelineOptions { Description(My option) Default.String(Hello world!) ValueProviderString getStringValue(); void setStringValue(ValueProviderString value); }要点自定义选项接口必须继承PipelineOptions并为每个选项声明getter / setter 对编程指南创建自定义选项一节的通用约定。getter 的返回类型声明为ValueProviderString这是让该选项成为运行时参数的关键。Description注解为该选项提供命令行帮助文本Default.String(Hello world!)声明默认值——当命令行未提供该参数时运行时get()会返回这个默认值。Java 中ValueProviderT支持的类型受限于简单类型集合若使用参数化ValueProviderTT须为受支持类型相关类型校验逻辑见 PipelineOptionsFactory.java。3.2 构建管道并添加日志分支public static void accessingValueProviderInfoAfterRunSnip1(String[] args) { MyOptions options PipelineOptionsFactory.fromArgs(args).withValidation().as(MyOptions.class); // Create pipeline. Pipeline p Pipeline.create(options); // Add a branch for logging the ValueProvider value. p.apply(Create.of(1)) .apply( ParDo.of( new DoFnInteger, Integer() { // Define the DoFn that logs the ValueProvider value. ProcessElement public void process(ProcessContext c) { MyOptions ops c.getPipelineOptions().as(MyOptions.class); // This example logs the ValueProvider value, but you could store it by // pushing it to an external database. LOG.info(Option StringValue was {}, ops.getStringValue()); } })); // The main pipeline. p.apply(Create.of(1, 2, 3, 4)).apply(Sum.integersGlobally()); p.run(); }逐步解读解析命令行参数PipelineOptionsFactory.fromArgs(args)按--optionvalue格式解析.withValidation()会检查必需参数并校验参数值.as(MyOptions.class)将通用选项对象投影为自定义接口类型工厂方法定义见 PipelineOptionsFactory.java。日志分支Create.of(1)产出含单个元素1的PCollectionIntegerParDo.of(...)应用日志DoFn。在 DoFn 内取回选项c.getPipelineOptions().as(MyOptions.class)从ProcessContext中取回当前作业的PipelineOptions再投影为MyOptions。这保证了读取的是该作业运行时的选项实例。记录值ops.getStringValue()返回的是ValueProviderString对象日志输出会隐式调用其toString()RuntimeValueProvider在可访问时toString()即返回get()的结果见 ValueProvider.java。如果你需要原值可显式调用.get()。主管道Create.of(1, 2, 3, 4)之后Sum.integersGlobally()求和代表真实业务逻辑与日志分支并行执行。运行示例# 不传自定义参数日志输出默认值 Hello world! java -cp ... org.apache.beam.examples.snippets.Snippets # 显式传入运行时参数日志输出 Option StringValue was my-value java -cp ... org.apache.beam.examples.snippets.Snippets --stringValuemy-value注意命令行选项名由 getter 方法名推导getStringValue→--stringValueBeam 会自动做驼峰到短横线的归一化处理。四、Python 实现完整代码与逐行讲解以下完整代码来自仓库示例 sdks/python/apache_beam/examples/snippets/snippets.py对应AccessingValueProviderInfoAfterRunSnip1片段。def accessing_valueprovider_info_after_run(): import logging import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions from apache_beam.options.value_provider import RuntimeValueProvider class MyOptions(PipelineOptions): classmethod def _add_argparse_args(cls, parser): parser.add_value_provider_argument(--string_value, typestr) class LogValueProvidersFn(beam.DoFn): def __init__(self, string_vp): self.string_vp string_vp # Define the DoFn that logs the ValueProvider value. # The DoFn is called when creating the pipeline branch. # This example logs the ValueProvider value, but # you could store it by pushing it to an external database. def process(self, an_int): logging.info(The string_value is %s % self.string_vp.get()) # Another option (where you dont need to pass the value at all) is: logging.info( The string value is %s % RuntimeValueProvider.get_value(string_value, str, )) beam_options PipelineOptions() args beam_options.view_as(MyOptions) # Create pipeline. with beam.Pipeline(optionsbeam_options) as pipeline: # Add a branch for logging the ValueProvider value. _ ( pipeline | beam.Create([None]) | LogValueProvs beam.ParDo(LogValueProvidersFn(args.string_value))) # The main pipeline. result_pc ( pipeline | main_pc beam.Create([1, 2, 3]) | beam.combiners.Sum.Globally())逐步解读自定义选项MyOptions(PipelineOptions)通过_add_argparse_args类方法向 argparse 注册选项。关键一步是调用parser.add_value_provider_argument(--string_value, typestr)而非普通的add_argument——该方法由_BeamArgumentParser提供pipeline_options.py内部会把type替换为_static_value_provider_of(value_type)包装函数、把default替换为一个RuntimeValueProvider实例从而让该参数具备运行时取值语义。日志 DoFnLogValueProvidersFn在__init__中持有ValueProvider引用构造时传入args.string_valueprocess方法中调用self.string_vp.get()取真实值。第二种取值方式RuntimeValueProvider.get_value(string_value, str, )是类方法级的全局取值value_provider.py无需在 DoFn 构造时传递引用通过选项名直接从运行时选项字典中读取取不到时回退到默认值。分支构建pipeline | beam.Create([None]) | LogValueProvs beam.ParDo(...)中Create([None])提供占位元素触发process命名步骤LogValueProvs便于在 DAG 可视化与监控中识别。主管道beam.Create([1, 2, 3]) | beam.combiners.Sum.Globally()求和使用with语句管理Pipeline生命周期退出上下文时自动run()。运行示例python -m apache_beam.examples.snippets.snippets # 日志输出The string_value is None / The string value is python -m apache_beam.examples.snippets.snippets --string_valuemy-value # 日志输出The string_value is my-value / The string value is my-value五、底层原理RuntimeValueProvider 与运行时上下文要真正理解为什么必须放进 DAG需要看ValueProvider的三种实现JavaValueProvider.javaPythonvalue_provider.py实现语义是否可在构建期取值StaticValueProvider值在构建期就已固定例如从命令行参数直接解析出的值是isAccessible()恒为trueRuntimeValueProvider值在作业执行期才可用例如 Dataflow 模板的运行时参数否运行上下文未就绪时抛异常NestedValueProvider包装另一个ValueProvider通过 translator 函数做值转换取决于被包装者5.1 Java 侧按 optionsId 查找运行时选项RuntimeValueProvider内部维护了一个全局注册表ConcurrentHashMapLong, PipelineOptions optionsMapValueProvider.java。作业启动时通过setRuntimeOptions将当前PipelineOptions按optionsId登记进去get()先从optionsMap按optionsId取出PipelineOptions若为null则抛出IllegalStateExceptionValue only available at runtime, but accessed from a non-runtime context这正是 DAG 外取值失败的原因随后通过反射调用对应的 getter 方法若反序列化得到的是StaticValueProvider说明作业执行期注入了新值则直接返回其值否则返回默认值。因此在 runner 的 worker 上执行DoFn.process时运行时选项必然已注册get()可以稳定取到值。5.2 Python 侧全局 runtime_optionsRuntimeValueProvider使用类属性runtime_options保存运行时选项value_provider.pyis_accessible()等价于RuntimeValueProvider.runtime_options is not Noneget()在runtime_options为None时抛出RuntimeValueProviderError作业运行期由set_runtime_options注入选项get_value类方法按选项名查值并做类型转换缺省回退默认值。这也解释了文档给出的模式为什么必须以 DoFn 处理占位值只有把取值动作放进 DAG 内的process方法才能确保它发生在运行时上下文已就绪之后。5.3 另一种官方途径Create.ofProvider除日志分支外若想把ValueProvider的值本身变成PCollection供下游变换消费Java SDK 还提供了Create.ofProvider(ValueProviderT, CoderT)Create.java。它会在运行时把 provider 的值作为一个元素注入管道——本质上是取值进 DAG的另一种封装与本文模式互补。六、模式扩展从记录日志到存储落地文档示例的DoFn注释明确给出了扩展方向This example logs the ValueProvider value, but you could store it by pushing it to an external database.在实际工程中你可以把process方法体内的日志调用替换为任意副作用操作例如将参数写入外部数据库 / 对象存储用于作业审计与事后追溯将参数追加到监控指标Beam Metrics中供面板展示根据参数值触发告警或后续补偿流程。实现时注意保持DoFn的幂等性runner 可能重试处理元素并在写入外部系统时做好去重或使用幂等写入语义。此外该分支使用单元素占位输入整体开销极小不会对主管道性能产生实质影响但若你的 runner 开启了流水线优化建议给分支步骤明确命名如LogValueProvs避免被合并或修剪。七、小结何时使用该模式适用场景使用ValueProvider参数化的管道尤其是 Dataflow 模板类作业需要在作业执行后审计、记录或落盘其实际使用的运行时参数。核心套路Create占位元素 →ParDo日志 DoFn分支 主业务变换主分支→ 统一run()在process内通过c.getPipelineOptions().as(MyOptions.class)Java或self.string_vp.get()/RuntimeValueProvider.get_value(...)Python读取参数。配套文档管道构建与选项配置的完整规范见 创建管道 与 配置管道选项变换PTransform应用方式见 应用变换。这套模式在 Apache Beam 官方文档中被列为 Pipeline Options 的核心 pattern双语言示例均可在本仓库中直接找到并运行Java 版位于 Snippets.javaPython 版位于 snippets.py可作为你接入自己业务的起点模板。赞分享批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载相关推荐Apache Beam 管道选项模式实战使用 ValueProvider 在作业运行后追溯记录运行时参数Apache Beam 管道选项模式实战使用 ValueProvider 在作业运行后追溯记录运行时参数 Apache Beam 允许通过 PipelineO大数据批处理流处理数据工程Apache Beam SQL 扩展使用 SET 与 RESET 语句配置 Pipeline OptionsApache Beam SQL 扩展使用 SET 与 RESET 语句配置 Pipeline Options Beam SQL 在标准 SQL 之外提供了一套大数据批处理流处理数据工程Apache Beam Hazelcast Jet Runner 实战指南本地与远程集群运行 WordCount 及全部 Pipeline 参数解析Apache Beam Hazelcast Jet Runner 实战指南本地与远程集群运行 WordCount 及全部 Pipeline 参数解析 Apac大数据批处理流处理数据工程上一篇【亲测免费】 PDFKit: 使用Ruby将HTMLCSS转化为PDF的利器下一篇2025最强WPF控件库WPFDevelopers从入门到精通创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考