PyFlink Table API 指标系统(Metrics)实战指南:从注册到上报的完整解析
大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载导读在 PyFlink 中编写 Python UDF 时如何观测自定义函数内部的运行状态例如处理了多少条数据、当前缓冲长度、每秒事件吞吐答案就是 PyFlink 提供的指标Metrics系统。本指南基于 Flink 仓库中 PyFlink Table API Metrics 官方文档系统讲解如何在 Python UDF 的open方法中通过function_context.get_metric_group()注册Counter、Gauge、Distribution、Meter四类指标介绍MetricGroup的作用域Scope与用户变量User Variables机制并结合仓库源码剖析其底层实现最后说明如何将指标对接 Reporter、REST API 与 Dashboard帮助读者完整掌握 PyFlink 指标从定义到暴露的整条链路。一、PyFlink 指标系统概述PyFlink 继承了 Flink 的指标体系一套允许用户采集并对外暴露指标gathering and exposing metrics to external systems的系统。在 Table API 的 Python UDF 中你可以像 Java 端RichFunction通过getRuntimeContext().getMetricGroup()获取指标组一样在 Python 端通过FunctionContext获得对指标系统的访问入口。核心入口如下def open(self, function_context): metric_group function_context.get_metric_group()其中function_context是FunctionContext类型的对象。在仓库源码 flink-python/pyflink/table/udf.py 中可以看到FunctionContext.get_metric_group()会返回当前并行子任务parallel subtask的MetricGroup如果指标功能未启用_base_metric_group为None它会抛出RuntimeError并提示通过python.metric.enabled配置开启。FunctionContext还额外提供了get_job_parameter(key, default_value)用于在 UDF 中读取全局作业参数自 1.17 版本起。注意open方法在 UDF 实际被调用前执行一次适合做注册指标、建立连接等一次性初始化工作详见 flink-python/pyflink/table/udf.py 中UserDefinedFunction.open的注释说明。指标开关python.metric.enabled指标是否可用由配置项python.metric.enabled控制其定义位于 flink-python/src/main/java/org/apache/flink/python/PythonOptions.javapublic static final ConfigOptionBoolean PYTHON_METRIC_ENABLED ConfigOptions.key(python.metric.enabled) .booleanType() .defaultValue(true) .withDescription( When it is false, metric for Python will be disabled. ...);该配置默认值为true即默认开启。在 Java 侧的算子实现中AbstractPythonFunctionOperator.getFlinkMetricContainer()会根据该配置决定是否创建FlinkMetricContainer见 flink-python/src/main/java/org/apache/flink/streaming/api/operators/python/AbstractPythonFunctionOperator.java测试用例 flink-python/src/test/java/org/apache/flink/python/PythonOptionsTest.java 也验证了默认值行为。若需显式关闭可在 TableEnvironment 配置中设置t_env.get_config().set(python.metric.enabled, false)关闭后调用get_metric_group()会抛出RuntimeError可参考 flink-python/pyflink/table/tests/test_udf.py 中测试“metric disabled”场景的写法。二、在 Python UDF 中注册指标PyFlink 支持四种指标类型Counter计数器、Gauge仪表、Distribution分布、Meter仪表/吞吐率。它们都通过MetricGroup上的对应方法注册抽象接口定义在 flink-python/pyflink/metrics/metricbase.py自 1.11.0 版本引入。下表概括了四种指标的核心特征指标类型注册方法更新方法语义值类型限制Countercounter(name: str)inc(n)/dec(n)计数可增可减整数Gaugegauge(name: str, obj: Callable[[], int])按需回调取值按需上报当前值仅整数Distributiondistribution(name: str)update(n: int)上报 sum/count/min/max/mean 分布信息仅整数Metermeter(name: str, time_span_in_seconds: int 60)mark_event(n)平均吞吐率事件数整数2.1 Counter计数器Counter用于计数。当前值可通过inc()/inc(n: int)增加或dec()/dec(n: int)减少。注册方式为在MetricGroup上调用counter(name: str)from pyflink.table.udf import ScalarFunction class MyUDF(ScalarFunction): def __init__(self): self.counter None def open(self, function_context): self.counter function_context.get_metric_group().counter(my_counter) def eval(self, i): self.counter.inc(i) return i上面示例中每处理一条数据就把i累加到计数器上从而统计 UDF 累计处理的值总和。Counter接口还提供get_count()返回当前计数值见 metricbase.py 中Counter抽象类。2.2 Gauge仪表Gauge按需提供当前值provides a value on demand。注册时传入一个可调用对象CallableFlink 在采样时调用该对象获取值。注册方式为gauge(name: str, obj: Callable[[], int])。PyFlink 的 Gauge 仅支持整数类型的值与 Java 端 Gauge 可返回任意类型不同from pyflink.table.udf import ScalarFunction class MyUDF(ScalarFunction): def __init__(self): self.length 0 def open(self, function_context): function_context.get_metric_group().gauge(my_gauge, lambda : self.length) def eval(self, i): self.length i return i - 1上例通过闭包捕获 UDF 实例的self.length属性使 Gauge 实时反映最近一次处理的输入值。在嵌入式embedded执行模式下Python 的 Gauge 回调会被包装成 Java 侧的org.apache.flink.python.metric.embedded.MetricGauge通过PythonGaugeCallable.get_value()调用 Python 函数取值见 flink-python/pyflink/fn_execution/metrics/embedded/metric_impl.py。2.3 Distribution分布Distribution报告已上报值的分布信息sum总和、count次数、min最小值、max最大值和 mean均值。通过update(n: int)更新数值通过distribution(name: str)注册。同样仅支持整数分布from pyflink.table.udf import ScalarFunction class MyUDF(ScalarFunction): def __init__(self): self.distribution None def open(self, function_context): self.distribution function_context.get_metric_group().distribution(my_distribution) def eval(self, i): self.distribution.update(i) return i - 1从测试 flink-python/pyflink/fn_execution/metrics/tests/test_metric.py 可以看到 Distribution 的聚合语义依次update(10)、update(2)后累计结果为DistributionData(12, 2, 2, 10)即 sum12、count2、min2、max10。该测试还展示了 Counter、Meter、Distribution 在 Process 模式基于 Apache Beam 的 metric 容器下的完整行为是理解指标底层聚合的很好参考。2.4 Meter吞吐率仪表Meter测量平均吞吐率average throughput。单次事件用mark_event()记录同时多次事件用mark_event(n: int)记录。注册方式为meter(name: str, time_span_in_seconds: int 60)其中time_span_in_seconds是计算平均速率的时间窗口跨度默认值为 60 秒。示例from pyflink.table.udf import ScalarFunction class MyUDF(ScalarFunction): def __init__(self): self.meter None def open(self, function_context): # 以 120 秒为窗口统计每秒平均事件数默认窗口为 60 秒 self.meter function_context.get_metric_group().meter(my_meter, time_span_in_seconds120) def eval(self, i): self.meter.mark_event(i) return i - 1从实现上看Process 模式下由于 Beam 没有原生 Meter 类型GenericMetricGroup.meter用Metrics.counter实现 Meter并将time_span_in_seconds拼入命名空间见 flink-python/pyflink/fn_execution/metrics/process/metric_impl.py而嵌入式Embedded模式下则直接使用 Java 侧的org.apache.flink.metrics.MeterView构造 Meter见 embedded/metric_impl.py。三、MetricGroup 的作用域Scope每个指标都会被赋予一个标识符identifier和一组键值对用来确定该指标在外部系统中的上报位置。PyFlink 中作用域的定义规则与 Java 侧一致详见 Flink 指标作用域定义文档。默认标识符分隔符为.可通过metrics.scope.delimiter配置修改。3.1 用户作用域User Scope通过MetricGroup.add_group(key: str, value: str None)定义用户作用域当value为None时创建一个普通子组generic sub-group该组被加入当前组的子组列表并返回新组当value不为None时创建一组key-value 形式的 MetricGroup 对key 组加入当前组的子组value 组加入 key 组的子组此时返回 value 组同时定义一个用户变量user variable。示例function_context \ .get_metric_group() \ .add_group(my_metrics) \ .counter(my_counter) function_context \ .get_metric_group() \ .add_group(my_metrics_key, my_metrics_value) \ .counter(my_counter)在 flink-python/pyflink/fn_execution/metrics/tests/test_metric.py 中test_add_group与test_add_group_with_variable分别验证了两种调用路径add_group(my_group)生成路径root.my_group而add_group(key, value)生成路径root.key.value与文档描述完全一致。3.2 系统作用域System Scope系统作用域由 Flink 根据作业、算子、子任务等信息自动生成PyFlink 不做特殊处理规则与 Java 侧完全一致。详细定义见 Flink 系统作用域文档。3.3 全部变量列表List of all Variables作用域格式中可以引用一组系统预定义变量如job_id、task_attempt_num、operator_name等。完整清单见 Flink 全部变量列表文档。3.4 用户变量User Variables通过MetricGroup.addGroup(key: str, value: str)并指定value参数即可定义用户变量例如function_context \ .get_metric_group() \ .add_group(my_metrics_key, my_metrics_value) \ .counter(my_counter)重要限制用户变量不能用于作用域格式scope formats中即不能用用户变量拼接指标标识符但可以用它做分组或过滤维度。从实现看Process 模式下GenericMetricGroup.add_group对(name, extra)的处理是先创建 key 类型子组再在其下创建 value 类型子组并返回后者见 process/metric_impl.py_add_group还会做去重——同名的同类型子组不会重复创建见同文件_add_group方法。命名空间通过 JSON 序列化的组名与组类型列表生成_get_namespace方法这解释了测试中断言的[my_group, MetricGroupType.generic]格式。四、与 Flink 共用的指标能力PyFlink 指标体系与 Flink 共用以下能力详见对应文档Reporter上报器指标最终由各 Reporter 定期上报到外部系统如 JMX、Graphite、Prometheus 等见 指标上报器文档。通用属性通过metrics.reporter.reporter_name.property配置例如metrics.reporters: my_jmx_reporter,my_other_reporter metrics.reporter.my_jmx_reporter.factory.class: org.apache.flink.metrics.jmx.JMXReporterFactory metrics.reporter.my_jmx_reporter.port: 9020-9040 metrics.reporter.my_jmx_reporter.scope.variables.excludes: job_id;task_attempt_num metrics.reporter.my_jmx_reporter.scope.variables.additional: cluster_name:my_test_cluster,tag_name:tag_value metrics.reporter.my_other_reporter.factory.class: org.apache.flink.metrics.graphite.GraphiteReporterFactory metrics.reporter.my_other_reporter.host: 192.168.1.1 metrics.reporter.my_other_reporter.port: 10000可以看到通过scope.variables.excludes和scope.variables.additional可以灵活裁剪或扩充上报时的作用域变量。系统指标System metricsFlink 自动采集的 CPU、内存、GC、网络等系统级指标见 系统指标文档延迟追踪Latency tracking追踪记录处理延迟见 延迟追踪文档REST API 集成通过 REST API 查询指标见 REST API 集成文档Dashboard 集成在 Flink Web UI 中可视化指标见 Dashboard 集成文档。这些能力对 PyFlink 用户透明可用——你在 Python UDF 中注册的指标会与 Java 算子指标一样进入同一套 Flink 指标体系通过上述通道对外暴露。五、两种执行模式下的底层实现PyFlink 的 Python UDF 支持两种执行模式指标系统的底层实现也因此分为两套理解这一点有助于排查指标异常5.1 Process 模式默认python.execution-mode默认值为process见 PythonOptions.java 附近。此模式下指标基于 Apache Beam 的 metrics 容器实现核心类是GenericMetricGroupflink-python/pyflink/fn_execution/metrics/process/metric_impl.pycounter→Metrics.counter(namespace, name)gauge→Metrics.gauge(namespace, name)同时将 Python 回调存入_flink_gaugemeter→ 由于 Beam 无 Meter 类型用Metrics.counter模拟time_span_in_seconds进入命名空间distribution→Metrics.distribution(namespace, name)。5.2 Embedded嵌入式模式嵌入式模式通过pemja直接调用 Java 侧 API核心类是MetricGroupImplflink-python/pyflink/fn_execution/metrics/embedded/metric_impl.pyadd_group→self._metrics.addGroup(name)或addGroup(name, extra)gauge→ 包装为 Java 类MetricGaugePythonGaugeCallable负责回调 Python 函数meter→ Java 类org.apache.flink.metrics.MeterViewtime_span_in_seconds作为窗口参数distribution→ Java 类org.apache.flink.python.metric.embedded.MetricDistribution通过 Gauge 形式注册。两种模式在 Java 算子侧都由PYTHON_METRIC_ENABLED配置控制是否构建FlinkMetricContainer相关链路见 AbstractPythonFunctionOperator.java 以及 Table 侧的 AbstractPythonScalarFunctionOperator.java 等算子实现。六、完整示例在 Table API 作业中使用指标将以上知识点串起来一个完整的 PyFlink Table API 指标使用流程如下from pyflink.table import EnvironmentSettings, TableEnvironment from pyflink.table.udf import ScalarFunction class MyUDF(ScalarFunction): def __init__(self): self.counter None self.meter None def open(self, function_context): metric_group function_context.get_metric_group() self.counter metric_group.counter(my_counter) # 120 秒窗口的平均吞吐率 self.meter metric_group.meter(my_meter, time_span_in_seconds120) # 用户作用域 用户变量示例 metric_group.add_group(my_metrics_key, my_metrics_value).counter(my_scoped_counter) def eval(self, i): self.counter.inc(1) self.meter.mark_event(1) return i env_settings EnvironmentSettings.in_streaming_mode() t_env TableEnvironment.create(env_settings) # 指标默认开启如需显式开启可设置 # t_env.get_config().set(python.metric.enabled, true) t_env.create_temporary_system_function(my_udf, MyUDF()) # 将 UDF 注册到查询中使用作业运行后即可在 Web UI / REST API / Reporter 中看到指标 result t_env.sql_query(SELECT my_udf(id) FROM my_source)运行作业后可以在 Flink Web UI 的指标面板或对接的 Reporter如 Prometheus、JMX中观察my_counter、my_meter、my_scoped_counter等指标的变化。七、常见问题与排查建议调用get_metric_group()抛出RuntimeError(Metric has not been enabled...)说明python.metric.enabled被显式关闭默认是开启的检查 TableEnvironment 配置或集群配置中是否设置了python.metric.enabledfalse。Gauge/Distribution 出现非预期值PyFlink 的 Gauge 与 Distribution仅支持整数传入非整数会与类型约定不符同时 Gauge 的值是“按需回调”要确保闭包捕获的引用在eval中被正确更新。指标标识符与预期不符检查作用域格式与用户变量使用。注意用户变量不能用于作用域格式若需要参与标识符拼接应改用系统变量。无法在 Dashboard 中看到自定义指标确认已为作业配置了至少一个有效的指标 Reporter见 metric_reporters.md并检查metrics.scope.delimiter等作用域配置是否影响了指标名解析。参考文档与源码索引本文主体PyFlink Table API Metrics 官方文档Flink 指标总览docs/content/docs/ops/metrics.md指标上报器docs/content/docs/deployment/metric_reporters.md指标抽象接口flink-python/pyflink/metrics/metricbase.pyFunctionContext与 UDF 基类flink-python/pyflink/table/udf.pyProcess 模式实现flink-python/pyflink/fn_execution/metrics/process/metric_impl.pyEmbedded 模式实现flink-python/pyflink/fn_execution/metrics/embedded/metric_impl.py指标配置项定义flink-python/src/main/java/org/apache/flink/python/PythonOptions.java指标行为测试flink-python/pyflink/fn_execution/metrics/tests/test_metric.py赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐PyFlink Table API 指标Metrics实战指南在 Python UDF 中注册与暴露自定义指标PyFlink Table API 指标Metrics实战指南在 Python UDF 中注册与暴露自定义指标 导读 本文围绕 Flink 仓库中 PyF大数据流处理批处理数据工程PyFlink Table API 行级操作实战Map / FlatMap / Aggregate / FlatAggregate 完整指南PyFlink Table API 行级操作实战Map / FlatMap / Aggregate / FlatAggregate 完整指南 PyFlink大数据流处理批处理数据工程Apache Dubbo Metrics指标详解监控服务健康度的关键指标Apache Dubbo Metrics指标详解监控服务健康度的关键指标 你是否曾因服务响应缓慢却找不到根源而困扰是否在排查分布式系统问题时缺乏有效的数据支RPC框架微服务后端服务注册发现上一篇从论文到实践Pythia-Intervention-70m-Deduped的训练干预方法论全解析下一篇Questgen.ai 项目常见问题解决方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考