新闻详情

新闻详情

首页 / 资讯中心 / 详情

Apache Beam Metrics 详解:在 Python 管道中定义、查询与导出运行指标

发布时间:2026/9/29 2:30:29来源:尧图网络
Apache Beam Metrics 详解:在 Python 管道中定义、查询与导出运行指标
大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Apache Beam 的 Metrics API 为批处理和流式管道提供了一套统一的运行时观测手段你可以在任意 DoFn 或 Transform 中声明命名指标在管道执行期间动态累积数据并在作业结束后甚至运行中按步骤、按命名空间查询结果。本文以 10_basic_metrics.md 为核心骨架结合 Apache Beam Python SDK 源码metric.py、metricbase.py、execution.py与完整可运行示例 wordcount_with_metrics.py带你从声明、更新、查询到外部导出Graphite、REST HTTP完整掌握 Beam 指标系统的使用与底层实现。Beam 指标是什么在 Apache Beam 编程模型中指标Metrics用于洞察管道的当前状态包括管道执行期间的状态。它们与日志、Counter 累加器等传统手段的关键区别在于指标是结构化、命名且可查询的可以由 Runner 统一收集、聚合与上报。Beam 指标有三个核心语义命名与作用域Named and scoped每个指标都有一个由namespace命名空间和name名称组成的标识并被作用域绑定到管道中的某个具体步骤step。查询时既可以按名称过滤也可以按步骤路径过滤。动态创建Created dynamically指标可以在管道执行过程中被动态创建不需要在定义阶段预先注册全部指标。降级容忍Fallback behavior如果某个 Runner 不支持部分指标上报功能其回退行为是丢弃该指标的更新而不是使整个管道失败。这一设计保证了指标系统的可用性不会成为作业运行的风险点。在 Python SDK 中用户面向的入口统一为apache_beam.metrics.Metrics类具体实现分散在 metric.py用户 API、metricbase.py指标接口定义与 execution.py内部执行与收集。内置指标类型文档明确指出 Apache Beam 提供多种内置指标类型。Python SDK 中通过 metricbase.py 定义了以下接口指标类型操作接口语义源码位置Counter计数器inc(n1)/dec(n1)计数可增可减累积当前值与增量适合统计元素数量metricbase.pyDistribution分布update(value)收集变量的统计信息count、sum、min、max、mean用于刻画数据分布metricbase.pyGauge仪表set(value)跟踪变量的最新值适合监控水位、当前批大小等瞬态状态metricbase.py从当前仓库源码结构看Python SDK 还在此基础上扩展了更多指标形态可以推断属于较新版本新增能力StringSetadd(value)报告执行期间出现的唯一字符串集合见 metricbase.pyBoundedTrieadd(value)以有界前缀树形式报告唯一的字符串集合见 metricbase.pyHistogramupdate(value)结合BucketType分桶统计百分位分布见 metricbase.py。其中Counter 与 Distribution 只接受整数Gauge 同样被限制为整数值这一点在Metrics.distribution与Metrics.gauge的 docstring 中均有明确说明metric.py。声明一个指标Metrics 类要声明指标统一使用beam.metrics.Metrics类。文档给出的声明片段如下self.words_counter Metrics.counter(self.__class__, words) self.word_lengths_counter Metrics.counter(self.__class__, word_lengths) self.word_lengths_dist Metrics.distribution(self.__class__, word_len_dist) self.empty_line_counter Metrics.counter(self.__class__, empty_lines)这是 wordcount_with_metrics.py 中WordExtractingDoFn.__init__的真实代码。其背后对应Metrics类的四个静态工厂方法metric.py工厂方法签名返回值Metrics.countercounter(namespace, name)DelegatingCounterMetrics.distributiondistribution(namespace, name, process_wideFalse)DelegatingDistributionMetrics.gaugegauge(namespace, name, process_wideFalse)DelegatingGaugeMetrics.string_setstring_set(namespace, name)DelegatingStringSetMetrics.bounded_triebounded_trie(namespace, name)DelegatingBoundedTrieMetrics.histogramhistogram(namespace, name, bucket_type)DelegatingHistogramnamespace 与 name 的规则namespace既可以传类也可以传字符串。Metrics.get_namespacemetric.py对类型参数的处理方式是拼接模块名.类名if isinstance(namespace, type): return {}.{}.format(namespace.__module__, namespace.__name__) elif isinstance(namespace, str): return namespace else: raise ValueError(Unknown namespace type)这就是为什么文档示例中直接传self.__class__——它会被展开成apache_beam.examples.wordcount_with_metrics.WordExtractingDoFn这样的完整命名空间。命名空间用于对相关指标分组并防止不同模块出现同名指标冲突MetricName的定义也印证了这一点metricbase.py“The namespace allows grouping related metrics together and also prevents collisions between multiple metrics of the same name”。同时MetricName要求用户指标必须同时提供非空 namespace 与非空 name否则抛出ValueErrormetricbase.py。在 DoFn 中更新指标DelegatingCounter、DelegatingDistribution、DelegatingGauge等委托类把更新动作委托给内部MetricUpdatermetric.py用户只需在process中调用对应方法def process(self, element): text_line element.strip() if not text_line: self.empty_line_counter.inc(1) # 空行计数 1 words re.findall(r[\w\], text_line, re.UNICODE) for w in words: self.words_counter.inc() # 单词计数 1 self.word_lengths_counter.inc(len(w)) # 字母长度计数可多次调用 self.word_lengths_dist.update(len(w)) # 词长分布更新 return words对应 wordcount_with_metrics.py。可以看到 Counter 支持在同一元素上多次inc如对每个单词累计长度Distribution 则通过update累积每次观测值。完整示例带指标的 WordCount文档提到“关于实现细节参见 Apache Beam GitHub 仓库中的 WordCount with metrics 示例”在本地仓库中即 wordcount_with_metrics.py。它是 WordCount 的完整可运行版本管道部分与标准 WordCount 一致lines p | read ReadFromText(known_args.input) counts ( lines | split (beam.ParDo(WordExtractingDoFn()).with_output_types(str)) | pair_with_one beam.Map(lambda x: (x, 1)) | group beam.GroupByKey() | count beam.Map(count_ones)) output counts | format beam.Map(format_result) output | write WriteToText(known_args.output) result p.run() result.wait_until_finish()运行方式该文件头部注释中注明了 beam-playground 配置python -m apache_beam.examples.wordcount_with_metrics --output output.txt--input默认指向gs://dataflow-samples/shakespeare/kinglear.txt可替换为本地文件--output为必填参数。示例的独特之处在于main()末尾的指标查询逻辑它展示了如何从运行结果中把指标取出来并写入日志。查询指标PipelineResult.metrics().query()指标定义在管道中查询则通过运行结果的metrics()接口完成。Beam Python 中查询的入口是result.metrics().query(filter)并配合MetricsFilter做名称过滤。from apache_beam.metrics.metric import MetricsFilter # 按名称过滤出空行计数器 empty_lines_filter MetricsFilter().with_name(empty_lines) query_result result.metrics().query(empty_lines_filter) if query_result[counters]: empty_lines_counter query_result[counters][0] logging.info(number of empty lines: %d, empty_lines_counter.result) # 按名称过滤出词长分布 word_lengths_filter MetricsFilter().with_name(word_len_dist) query_result result.metrics().query(word_lengths_filter) if query_result[distributions]: word_lengths_dist query_result[distributions][0] logging.info(average word length: %d, word_lengths_dist.result.mean)对应 wordcount_with_metrics.py。有几个要点查询结果按类型分组query()返回的 dict 包含counters、distributions、gauges以及源码中定义string_sets、bounded_tries、histograms等键见 metric.py 的MetricResults常量过滤语义MetricsFilter().with_name(name)只匹配指标名称还支持with_namespace(...)与with_step(...)。底层匹配逻辑在 metric.py名称匹配要求 namespace 与 name 同时命中过滤条件步骤匹配则采用/分隔路径的连续子路径判断_matches_sub_path因此可以用步骤前缀或中间片段过滤result属性的回退语义MetricResult.result在committed未填充时自动回退到attempted这正是文档所说“Runner 不支持某些上报时降级”在查询层的体现execution.pyDistribution 的结果对象支持mean等统计量此处word_lengths_dist.result.mean直接给出平均词长。另外示例中查询前有一句注释与判断“Do not query metrics when creating a template which doesnt run”即通过hasattr(result, has_job)判断当前 Runner 是否真的执行了作业Direct Runner 没有该属性避免在只创建模板未运行时查询空指标。底层实现MetricUpdater、MetricCell 与 MetricsContainer要理解指标“为什么能跨线程、跨 bundle 累积”可以看内部实现仅供内部使用不保证向后兼容1.MetricUpdaterexecution.py所有委托指标DelegatingCounter 等的更新最终都落到这个可调用对象上。它把(cell_type, metric_name)组合成_TypedMetricName作为键调用时若process_wideTrue直接更新进程级容器PROCESS_WIDE_METRICS_CONTAINER否则通过statesampler.get_current_tracker()拿到当前状态采样器把更新挂到当前执行的 bundle 上下文中tracker.update_metric(...)。2.MetricCell家族cells.py每种指标类型对应一个 cell例如CounterCell维护一个整数value更新时持有线程锁self._lock以保证并发安全DistributionCell内部维护DistributionData累积 count/sum/min/max。每个 bundle、每个步骤都有独立的 cell最终由 Runner 聚合。3.MetricsContainerexecution.py按“步骤 bundle”持有该上下文内所有指标 cellself.metricsdict键为_TypedMetricName。get_cumulative()按类型把 cell 的累积值导出为MetricUpdatesto_runner_api_monitoring_infos(transform_id)则把每个 cell 序列化为 Runner 可消费的MonitoringInfocells.py 中还会记录指标首次更新的 UTC 起始时间。4.MetricResult的 committed / attemptedexecution.py物理更新attempted是执行期间已发生但未必提交的更新逻辑更新committed是已提交的更新。查询结果同时暴露两者result属性优先取 committed、缺省时回退 attempted。指标导出到外部 Sink文档明确说明指标可以导出到外部 SinkSpark 和 Flink Runner 支持 REST HTTP 与 Graphite。仓库中能找到对应的实现证据Graphite SinkMetricsGraphiteSink[runners/extensions-java/metrics/src/main/java/org/apache/beam/runners/extensions/metrics/MetricsGraphiteSink.java及其测试 MetricsGraphiteSinkTest.javaSpark Runner 侧也有CodahaleGraphiteSink[runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/sink/CodahaleGraphiteSink.java与GraphiteSink[runners/spark/src/main/java/org/apache/beam/runners/spark/metrics/sink/GraphiteSink.javaREST HTTP 上报由 Runner 侧配置将收集到的指标通过 HTTP 推送到外部监控服务具体开关与地址以所部署 RunnerSpark/Flink的配置文档为准。也就是说Beam 的指标体系分为两段SDK 侧负责定义、累积并序列化MonitoringInfoRunner 侧负责聚合并上报到外部系统Graphite、REST HTTP 等。如果你的 Runner 不支持某项上报指标更新会被静默丢弃而不会拖垮管道——这正是第一节提到的降级语义。小结声明Metrics.counter / distribution / gauge外加仓库中的string_set、bounded_trie、histogramnamespace 传类或字符串指标绑定到具体步骤更新在 DoFn 的process中调用inc、update、set内部由MetricUpdater路由到对应 bundle 的MetricCell线程安全地累积查询result.metrics().query(MetricsFilter().with_name(...))结果按counters / distributions / gauges等类型分组result属性自动回退 attempted导出Spark、Flink Runner 支持 Graphite 与 REST HTTP 外部 Sink可靠性Runner 不支持时丢弃指标更新而非失败保证指标系统的非侵入性。建议进一步阅读仓库内 metric_test.py 与 execution_test.py 了解各指标类型的边界行为或直接运行python -m apache_beam.examples.wordcount_with_metrics --output output.txt观察控制台输出的number of empty lines与average word length日志。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam Metrics监控指标详解自定义指标开发Apache Beam Metrics监控指标详解自定义指标开发 在数据处理管道中监控指标Metrics是评估系统性能和业务状态的关键工具。Apache批处理流处理大数据Apache Beam Go SDK 管道单元测试实战passert 断言与 ptest 运行器详解Apache Beam Go SDK 管道单元测试实战passert 断言与 ptest 运行器详解 Apache Beam 采用用户代码构建管道图、由分布大数据批处理流处理数据工程如何用Apache Beam监控生产管道Metrics指标与任务调试完整指南如何用Apache Beam监控生产管道Metrics指标与任务调试完整指南 Apache Beam 是统一的批流一体数据处理编程模型而 监控生产管道 的可大数据批处理流处理数据工程上一篇描述下一篇AirLLM基准测试与其他推理框架的性能对比分析创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
网站建设高端定制企业官网
RELATED

相关资讯

更多精彩内容,欢迎继续阅读

较早相关资讯

最新相关资讯

深入剖析Qt QFont:字体匹配、HiDPI与多语言填坑指南 2026/9/29 3:24:46

深入剖析Qt QFont:字体匹配、HiDPI与多语言填坑指南

写Qt界面的人,迟早都会被字体问题缠上。不是文字显示成方块,就是不同分辨率下控件错位,再不就是高分屏上字体发虚。这些问题绕来绕去,最后都会落到同一个类上——QFont。这个类表面上就是“设置字体名字和大小”,但它背…

阅读更多 →
用AI写代码后,为什么我们反而更累了?TaoToken配置排查与验证指南 2026/9/29 3:24:46

用AI写代码后,为什么我们反而更累了?TaoToken配置排查与验证指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →
ROS Noetic入门:从catkin工作空间到话题通信实战 2026/9/29 3:24:46

ROS Noetic入门:从catkin工作空间到话题通信实战

做机器人开发,尤其是做移动底盘、机械臂、无人车这类项目,绕过 ROS 基本是不可能的。哪怕你只是想把一个摄像头数据用起来,或者让电机按照速度指令转起来,最终都会落到同一个问题上:怎么把各个模块的程序放到同一个“消…

阅读更多 →
YOLOv8 目标检测全解析:网络结构、训练调参到部署实践 2026/9/29 3:24:33

YOLOv8 目标检测全解析:网络结构、训练调参到部署实践

上周有个做安防的哥们找我,说他们那套跑了快两年的 YOLOv5 检测流水线,老板发话了,问要不要迁到 YOLOv8。他这句话我这一年听过不下二十遍,每次我都得先反问一句:你是想要那两三个点的 mAP 提升,还是想让训…

阅读更多 →
Model-Optimizer实战:量化剪枝与蒸馏的模型压缩优化流程 2026/9/29 3:24:33

Model-Optimizer实战:量化剪枝与蒸馏的模型压缩优化流程

1. 模型优化器到底在优化什么第一次听到 Model-Optimizer 这个词,很多人会下意识地把它和“训练加速器”画等号,觉得无非就是让模型跑得快一点。但真正在项目里用过一轮之后你会发现,它解决的核心问题其实是在有限算力和显存预算下&#xff0…

阅读更多 →
STM32CubeMX 6.14全流程实战:从安装配置到代码生成与避坑指南 2026/9/29 3:24:26

STM32CubeMX 6.14全流程实战:从安装配置到代码生成与避坑指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

阅读更多 →

今日资讯

本周资讯

本月资讯

看完文章仍有疑问?

联系尧图顾问,获取一对一建站咨询

立即免费咨询 📞 400-888-8888
📞 ✉