大数据流处理批处理数据工程【免费下载链接】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),仅供参考
企业数字化 ERP 产品动态
相关推荐
MATLAB模式识别实战:源码解析与工业应用 1. 模式识别与MATLAB的黄金组合模式识别作为人工智能领域的核心技术之一,已经渗透到我们生活的方方面面。从手机人脸解锁到医疗影像分析,从工业质检到金融风控,这项技术正在重塑各行各业的运作方式。而MATLAB作为工程计算领域的"瑞士军刀… · 2026/9/23 15:12:38
SoC低功耗手册中文化实践:UPF与电源域隔离策略 简介:Low Power Methodology Manual for System-on-Chip Design 的中文翻译文档,面向数字IC设计工程师、SoC架构师及低功耗方向初学者,系统梳理芯片功耗问题的来源、动态与静态功耗的组成,并给出Multi-Vt、Power Gating、VTCMOS、… · 2026/9/23 15:12:37
BrowserSkill:基于CDP的浏览器自动化命令行工具 1. 项目概述:BrowserSkill不是浏览器,而是浏览器的“外科手术刀”BrowserSkill——这个名字乍一听像某个新出的浏览器,或者Chrome/Edge的某个隐藏功能模块。但实际接触过的人会立刻意识到:它根本不是UI界面产品,而是一… · 2026/9/23 15:12:31
DCH01隔离电源模块拆解:1W DC/DC转换器如何实现3kV隔离与稳定供电 简介:TI DCH01系列1W微型DC/DC转换器技术资料(PDF),面向电源设计、工业电子及嵌入式系统工程师,用于了解具备3kV隔离能力的非稳压转换器选型与应用。资料重点介绍该款5V输入、可输出单路/双路多种电压的模块࿰… · 2026/9/23 15:57:43
ARIS 跨阶段发现日志实战:用 FINDINGS_TEMPLATE 沉淀研究洞察与工程经验 ARIS 跨阶段发现日志实战:用 FINDINGS_TEMPLATE 沉淀研究洞察与工程经验 【免费下载链接】Auto-claude-code-research-in-sleep ARIS ⚔️ (Auto-Research-In-Sleep) — Lightweight Markdown-only skills for autonomous ML research: cross-model review loops, i… · 2026/9/23 15:57:43
RobotGo 跨平台桌面自动化完全指南:环境依赖、无 Cgo 纯 Go 构建与实战示例 RobotGo 跨平台桌面自动化完全指南:环境依赖、无 Cgo 纯 Go 构建与实战示例 【免费下载链接】robotgo RobotGo, Go Native cross-platform RPA, GUI automation, Auto test and Computer use vcaesar 项目地址: https://gitcode.com/gh_mirrors/ro/robotgo 本… · 2026/9/23 15:57:43
搞定硬盘作用原理,3个高频面试题轻松过 搞定硬盘作用原理,3个高频面试题轻松过 官方文档翻了几页就头大?别慌。 想搞懂 硬盘作用 在存储链路里的真实角色? 这些 高频面试题 背后其实只有三层逻辑。 项目目标与痛点拆解… · 2026/9/23 15:57:43
3招搞定手机怎么下载微信面试难题实战项目解析 3招搞定手机怎么下载微信面试难题实战项目解析 面试被问“手机怎么下载微信”背后的原理,90%的人答不上来。别笑,这看似弱智的问题,实则是考察你对移动应用分发机制、安全校验及网络协议理解的试金石。我带过不少校招新人,他们背了八股文,却连一个A… · 2026/9/23 0:00:03
你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 面试被问“高并发下如何保证消息不丢失”,你张口就是“用Redis”,结果面试官追问“如果Redis宕机了怎么办”,你瞬间卡壳。这种场景太常见了,很多新手在背八股文时,只记住了技术名词… · 2026/9/23 0:00:29