- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
本篇指南聚焦 Apache Beam 编程模型中的管道指标(Pipeline Metrics),讲解其在数据管道运行期如何提供可观测性、内置的指标类型(Counter / Distribution / Gauge 等)、Python 与 Java SDK 的声明与查询方式,并结合本仓库中的源码实现(metric.py、cells.py)与官方示例(wordcount_with_metrics.py)剖析其底层工作原理,最后给出 Spark / Flink Runner 将指标导出到 Graphite 等外部监控系统的配置方法。读完本文,你将能够在自己的 Beam 管道中声明、采集、查询并导出运行期指标。
一、什么是 Apache Beam 的管道指标
在 Apache Beam 模型中,**指标(metrics)**为用户提供了关于管道当前状态的可观测性信息,包括管道执行过程中的实时运行状况。它们是诊断管道性能、监控业务计数(如处理了多少条记录、出现多少空行、元素长度分布如何)的基础设施。
指标具有以下核心特性:
- 有命名(named)且限定(scoped)到管道中的特定步骤(step):每个指标都由
namespace + name唯一标识,并归属于产生它的具体 PTransform 步骤,因此可以在查询时精确追溯「哪个算子、哪个指标、什么值」。 - 可在管道执行期间动态创建:指标不要求在编译期预先注册,DoFn 实例可以在运行时按需创建与使用指标对象。
- 失败时优雅降级:如果某个 Runner 不支持指标上报的某一部分能力,默认的降级行为是丢弃(drop)该指标的更新,而不是让整个管道失败。这一设计保证了可观测性功能永远不会成为管道稳定性的负担。
从指标的生命周期看,一次指标更新会经历「物理更新(attempted)→ 提交更新(committed)」两个阶段。在 execution.py 中,MetricResult同时记录committed(已提交的逻辑更新)与attempted(执行期间发生的物理更新)两份数据,而MetricResult.result属性会优先返回committed,若 Runner 未填充则回退到attempted。
二、Beam 内置指标类型
Beam 提供了一组开箱即用的内置指标类型,覆盖计数、分布统计与瞬时值三类最常见需求:
| 指标类型 | 语义 | 典型操作 | 可统计的结果 |
|---|---|---|---|
| Counter(计数器) | 单调增减的整型计数 | inc()/dec(),默认每次inc增 1 | 累计值(int) |
| Distribution(分布) | 记录一组整型样本的统计量 | update(value) | sum、count、min、max、mean |
| Gauge(仪表/瞬时值) | 跟踪变量在任意时刻的最新值 | set(value) | 最新值(value)+ 时间戳(timestamp) |
2.1 各类型的实现细节(Python SDK)
这些接口的抽象定义位于 metricbase.py:
Counter:inc(n=1)递增、dec(n=1)等价于inc(-n);Distribution:通过update(value)累积样本,其结果对象DistributionResult提供sum、count、min、max、mean五个属性(见 cells.py),其中mean为sum / count,当count == 0时返回None;Gauge:通过set(value)覆盖最新值。
值得注意的底层约束:Distribution 只接受整数样本,Gauge 只接受整数值。这一点在 metric.py 的 docstring 中明确标注("Distribution metrics are restricted to integer-only distributions"、"Gauge metrics are restricted to integer-only values"),并在 cells.py 与 cells.py 的GaugeData/DistributionData注释中再次确认。
2.2 更多内置类型:StringSet、BoundedTrie 与 Histogram
除了文档中列出的三类经典指标,本仓库的 Python SDK 还实现了更多类型,可作为深入学习的扩展点(它们同样通过Metrics工厂类创建,见 metric.py):
- StringSet:记录去重后的字符串集合,
add(value)添加元素;其数据类StringSetData对总容量设限——所有字符串长度之和不得超过 1 MB,超出后继续添加的元素会被丢弃并打印警告(cells.py); - BoundedTrie:以带界前缀树形式记录字符串序列,默认界为 100 个节点,超出时会对最小分支执行裁剪(cells.py),常用于血缘(lineage)追踪场景;
- Histogram:基于桶(bucket)的分位数统计,结果对象可读取
p90、p95、p99分位数(cells.py);目前其上报到 Runner 的方法返回None,属于 worker 本地、内部使用为主的能力(cells.py)。
三、声明指标:beam.metrics.Metrics工厂类
要声明一个指标,统一使用beam.metrics.Metrics工厂类(Python SDK 导入路径为apache_beam.metrics.Metrics,底层实现在 metric.py)。
3.1 在 DoFn 中声明与使用指标
以官方示例 wordcount_with_metrics.py 为例,在 DoFn 的构造函数中一次性声明全部指标:
from apache_beam.metrics import Metrics class WordExtractingDoFn(beam.DoFn): def __init__(self): super().__init__() 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') 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要点说明:
- namespace 传类对象:
self.__class__会被Metrics.get_namespace()转换为'{module}.{ClassName}'形式的字符串命名空间(metric.py),例如__main__.WordExtractingDoFn;你也可以直接传字符串作为 namespace。 - namespace + name 共同构成指标全名:
MetricName由namespace与name组成,namespace 用于将相关指标分组并防止同名冲突(metricbase.py)。 - 在构造函数中创建、在
process中更新:这是 Beam 官方示例推荐的标准模式,避免每条元素都重复创建指标对象带来的开销。
3.2 工厂类的方法签名与默认值
Python SDK 中Metrics提供以下静态工厂方法(metric.py):
| 方法 | 参数 | 返回类型 | 备注 |
|---|---|---|---|
Metrics.counter(namespace, name) | namespace: 类或字符串;name: 字符串 | DelegatingCounter | 每次inc()默认 +1 |
Metrics.distribution(namespace, name, process_wide=False) | 同上 | DelegatingDistribution | 仅整数样本 |
Metrics.gauge(namespace, name, process_wide=False) | 同上 | DelegatingGauge | 仅整数值 |
Metrics.string_set(namespace, name) | 同上 | DelegatingStringSet | 去重字符串集合 |
Metrics.bounded_trie(namespace, name) | 同上 | DelegatingBoundedTrie | 带界前缀树 |
Metrics.histogram(namespace, name, bucket_type, logger=None) | 额外需要BucketType | DelegatingHistogram | 分位数统计 |
其中process_wide参数控制指标的作用域:为False(默认)时指标按当前 bundle 统计;为True时按整个进程统计(metric.py)。
3.3 Java SDK 中的等价 API
Java SDK 提供对称的 API,位于 Metrics.java,同样支持字符串与Class<?>两种 namespace 形式:
import org.apache.beam.sdk.metrics.Counter; import org.apache.beam.sdk.metrics.Distribution; import org.apache.beam.sdk.metrics.Gauge; import org.apache.beam.sdk.metrics.Metrics; private final Counter emptyLines = Metrics.counter(WordExtractingDoFn.class, "empty_lines"); private final Distribution wordLengths = Metrics.distribution(WordExtractingDoFn.class, "word_len_dist"); private final Gauge backlog = Metrics.gauge(WordExtractingDoFn.class, "backlog");Java 侧对应的结果类型为 DistributionResult.java(create(sum, count, min, max)静态工厂)与 GaugeResult.java(create(value, Instant timestamp))。
四、查询指标:result.metrics().query()与 MetricsFilter
指标在管道运行结束后(或运行期间,取决于 Runner)可以通过PipelineResult.metrics().query()查询。
4.1 完整查询示例
继续沿用 wordcount_with_metrics.py 中的做法:
result = p.run() result.wait_until_finish() # 按指标名过滤,查询 Counter 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) # 按指标名过滤,查询 Distribution 并读取均值 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)几点实操注意事项:
- 查询结果是一个字典,键包括
'counters'、'distributions'、'gauges'(以及string_sets、bounded_tries、histograms),每个键对应的值为MetricResult列表(metric.py 中MetricResults.query()的 docstring 给出了返回结构示例)。 - 查询前先判空:示例中每次查询后都检查对应列表是否为空,这是防御性写法——当 Runner 不支持某类指标时对应列表可能为空。
- 模板场景跳过查询:示例代码特意判断
hasattr(result, 'has_job') or result.has_job,即"仅在管道真正运行(而非仅创建模板)时才查询指标"。
4.2 MetricsFilter 的过滤维度
MetricsFilter(定义于 metric.py)提供三类过滤条件,且只作用于用户自定义指标:
| 方法 | 作用 |
|---|---|
with_name('name')/with_names([...]) | 按指标名过滤(传入字符串会报错,必须传可迭代集合) |
with_namespace(cls_or_str)/with_namespaces([...]) | 按命名空间过滤,内部会自动把类转换为字符串 namespace |
with_step('step')/with_steps([...]) | 按步骤(PTransform 名称)过滤 |
匹配逻辑(metric.py)由各 Runner 实现:_matches_name同时校验 namespace 与 name;_matches_scope将过滤条件中的 step 按/拆分为子路径,只要它是实际 step 路径的连续子序列即视为匹配(_is_sub_list),支持模糊的步骤层级匹配。
五、指标如何工作:从 MetricCell 到 MonitoringInfo
理解指标上报链路,可以按「更新 → 汇聚 → 上报」三步来看,对应 Python SDK 中三个核心模块:
- 更新:MetricCell 累积内存变更。每个指标在「每个上下文(context)× 每个 bundle」中拥有一个独立的 cell,负责线程安全地累积变更(cells.py)。
CounterCell维护一个整型value,update()时加锁累加(Cython 编译场景下利用 GIL 免锁);DistributionCell维护DistributionData(count/sum/min/max 四元组),每次update更新计数与极值;GaugeCell直接覆盖value并记录time.time()时间戳(cells.py)。所有 cell 都实现了combine(),供 Runner 跨 bundle 聚合。 - 组织:MetricsContainer 与 MetricKey。
MetricsContainer持有单个步骤、单个提交单元(bundle)内的全部指标;MetricKey由 step 名 +MetricName(namespace+name)+ 附加 labels 唯一标识一个指标实例(execution.py)。 - 上报:序列化为 MonitoringInfo。每个 cell 的
to_runner_api_monitoring_info()会生成带start_time的MonitoringInfo(cells.py):用户 Counter 被编码为int64_user_counter、Distribution 为int64_user_distribution、Gauge 为int64_user_gauge等(cells.py),随后由 Runner 汇聚并对外报告。
这条链路解释了为什么指标"可以在执行期间动态创建"——创建只是向环境注册一个MetricUpdater,真正读写发生在 cell 层,与管道数据路径解耦;也解释了"Runner 不支持时丢弃更新"的降级语义——上报环节的失败不会回传影响数据处理。
六、将指标导出到外部系统(Spark / Flink Runner)
Beam 允许将指标导出到外部监控 sink。Spark 与 Flink Runner 支持通过 REST HTTP 与 Graphite 导出指标。
6.1 Spark Runner 导出到 Graphite
本仓库为 Spark Runner 提供了专用的 Graphite sink 实现:
- GraphiteSink.java(经典 Spark Runner);
- CodahaleGraphiteSink.java(Spark Structured Streaming 路径)。
两者都委托 Spark 自带的org.apache.spark.metrics.sink.GraphiteSink,通过 Spark metrics 配置启用,可上报包括 Beam step 指标在内的监控数据。典型配置(写入 Spark metrics 配置,如spark.metrics.conf):
spark.metrics.conf.*.sink.graphite.class=org.apache.beam.runners.spark.metrics.sink.GraphiteSink spark.metrics.conf.*.sink.graphite.host=<graphite_hostname> spark.metrics.conf.*.sink.graphite.port=<graphite_listening_port> spark.metrics.conf.*.sink.graphite.period=10 spark.metrics.conf.*.sink.graphite.unit=seconds spark.metrics.conf.*.sink.graphite.prefix=<optional_prefix> spark.metrics.conf.*.sink.graphite.regex=<optional_regex_to_send_matching_metrics>字段含义:host/port指定 Graphite 服务地址;period与unit控制上报周期(示例为每 10 秒);prefix可选,为指标名添加前缀;regex可选,仅上报匹配该正则的指标。Structured Streaming 变体只需将class换成org.apache.beam.runners.spark.structuredstreaming.metrics.sink.CodahaleGraphiteSink。
6.2 关于 REST HTTP 导出
Flink Runner 侧同样具备指标外部化能力(可通过 REST/HTTP 接口与 job 通信并获取运行指标),而 Spark 侧的 REST HTTP 导出基于其自带 metrics 系统的相应 sink 配置。具体配置参数以对应 Runner 发行版文档为准,建议在实际部署时结合监控后端(Graphite、Prometheus 等)的协议要求做选型。
七、动手实践:运行带指标的 WordCount
仓库中的 wordcount_with_metrics.py 是可直接运行的完整示例(同时具备 beam-playground 元数据标注,见文件头部注释)。运行方式:
# 在仓库根目录下,使用 Python SDK 本地运行(Direct Runner) python -m apache_beam.examples.wordcount_with_metrics --output output.txt默认输入为gs://dataflow-samples/shakespeare/kinglear.txt(也可用--input指定本地文件),运行结束后控制台会输出两类日志:
number of empty lines: N(来自empty_linesCounter);average word length: N(来自word_len_distDistribution 的mean)。
实践提示:
- 若要观察更细粒度的指标,可在查询阶段改用
MetricsFilter().with_step('split')等过滤器,缩小到特定 PTransform 步骤; - 指标值会随 Runner 的实现而呈现 committed/attempted 差异,读取时优先使用
MetricResult.result(自动回退); - 若某 Runner 不支持你使用的指标类型,管道不会失败,但该指标不会被报告——这是 Beam 指标"优雅降级"语义的直接体现。
八、总结
本文以 Beam 官方指标文档为核心,完整覆盖了:指标的三类核心类型(Counter / Distribution / Gauge)及其语义、beam.metrics.Metrics工厂类的声明方式(Python 与 Java 双 SDK)、result.metrics().query()与MetricsFilter的查询过滤方法、从 MetricCell 到 MonitoringInfo 的底层上报链路,以及 Spark / Flink Runner 向 Graphite / REST HTTP 导出指标的具体配置。相关可进一步研读的源码入口包括 metric.py、metricbase.py、cells.py、execution.py 与 Metrics.java。掌握了这些能力,你就可以为自己的 Beam 管道建立一套可查询、可导出、可监控的运行期可观测性体系。
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam Metrics 详解:在 Python 管道中定义、查询与导出运行指标
Apache Beam Metrics 详解:在 Python 管道中定义、查询与导出运行指标 Apache Beam 的 Metrics API 为批处理和流
大数据批处理流处理数据工程如何用Apache Beam监控生产管道?Metrics指标与任务调试完整指南
如何用Apache Beam监控生产管道?Metrics指标与任务调试完整指南 Apache Beam 是统一的批流一体数据处理编程模型,而 监控生产管道 的可
后端音视频前端HackRF 频谱分析仪快速上手实战:5 步把入门 SDR 调校成专业频谱监测台
HackRF 频谱分析仪快速上手实战:5 步把入门 SDR 调校成专业频谱监测台 深夜两点,值班工程师盯着 2.4GHz 频段屏幕上那条突然拱起的尖峰——这已经
后端RPC框架
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考