Flink Kubernetes Operator自动扩缩容缺失TM指标问题排查
Flink Kubernetes Operator自动扩缩容报错:缺失必要的TM指标解决方案
错误日志
java.lang.RuntimeException: Missing required TM metrics at org.apache.flink.autoscaler.RestApiMetricsCollector.queryTmMetrics(RestApiMetricsCollector.java:179) at org.apache.flink.autoscaler.ScalingMetricCollector.updateMetrics(ScalingMetricCollector.java:137) at org.apache.flink.autoscaler.JobAutoScalerImpl.runScalingLogic(JobAutoScalerImpl.java:183) at org.apache.flink.autoscaler.JobAutoScalerImpl.scale(JobAutoScalerImpl.java:103) at org.apache.flink.kubernetes.operator.reconciler.deployment.AbstractFlinkResourceReconciler.applyAutoscaler(AbstractFlinkResourceReconciler.java:219) at org.apache.flink.kubernetes.operator.reconciler.deployment.AbstractFlinkResourceReconciler.reconcile(AbstractFlinkResourceReconciler.java:142) at org.apache.flink.kubernetes.operator.controller.FlinkDeploymentController.reconcile(FlinkDeploymentController.java:155) at org.apache.flink.kubernetes.operator.controller.FlinkDeploymentController.reconcile(FlinkDeploymentController.java:62) at io.javaoperatorsdk.operator.processing.Controller$1.execute(Controller.java:153) at io.javaoperatorsdk.operator.processing.Controller$1.execute(Controller.java:111) at org.apache.flink.kubernetes.operator.metrics.OperatorJosdkMetrics.timeControllerExecution(OperatorJosdkMetrics.java:80) at io.javaoperatorsdk.operator.processing.Controller.reconcile(Controller.java:110) at io.javaoperatorsdk.operator.processing.event.ReconciliationDispatcher.reconcileExecution(ReconciliationDispatcher.java:136) at io.javaoperatorsdk.operator.processing.event.ReconciliationDispatcher.handleReconcile(ReconciliationDispatcher.java:117) at io.javaoperatorsdk.operator.processing.event.ReconciliationDispatcher.handleDispatch(ReconciliationDispatcher.java:91) at io.javaoperatorsdk.operator.processing.event.ReconciliationDispatcher.handleExecution(ReconciliationDispatcher.java:64) at io.javaoperatorsdk.operator.processing.event.EventProcessor$ReconcilerExecutor.run(EventProcessor.java:452) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(Unknown Source) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(Unknown Source) at java.base/java.lang.Thread.run(Unknown Source)
当前Flink配置
job.autoscaler.metrics.window: 3m taskmanager.memory.jvm-metaspace.size: 256 mb metrics.system-resource: 'true' pipeline.max-parallelism: '24' taskmanager.network.detailed-metrics: 'true' job.autoscaler.target.utilization.boundary: '0.1' job.autoscaler.catch-up.duration: 5m job.autoscaler.restart.time: 2m job.autoscaler.scaling.enabled: 'true' job.autoscaler.stabilization.interval: 1m job.autoscaler.enabled: 'true' jobmanager.scheduler: adaptive
解决方法
这个错误是因为Autoscaler无法从TaskManager获取到必需的运行指标,你可以从以下几个方向排查解决:
补充指标收集相关配置
在你的Flink配置中添加以下项,确保TaskManager的指标能被正确收集并通过JobManager的REST API暴露:# 开启指标采集工厂,确保TaskManager指标被正常收集 metrics.reporter.prometheus.factory.class: org.apache.flink.metrics.prometheus.PrometheusReporterFactory # 确保指标作用域符合默认规则,若之前有修改需恢复 metrics.scope.jm: "<host>.jobmanager" metrics.scope.tm: "<host>.taskmanager.<tm_id>" metrics.scope.task: "<host>.taskmanager.<tm_id>.<job_id>.<task_name>.<subtask_index>"验证JobManager REST API的可访问性
- 进入Operator的Pod,执行
curl <jobmanager-service-name>:8081/taskmanagers,检查是否能返回所有TaskManager的注册信息。 - 执行
curl <jobmanager-service-name>:8081/metrics,搜索是否包含taskmanager_job_task_numBytesInPerSecond、taskmanager_job_task_numBytesOutPerSecond这类TaskManager核心指标。 - 如果无法访问,检查JobManager的Service标签选择器是否匹配Pod,以及是否有NetworkPolicy阻止了Operator与JobManager的通信。
- 进入Operator的Pod,执行
确认版本兼容性
确保Flink Kubernetes Operator版本与Flink版本匹配,比如Operator 1.7.x对应Flink 1.17+,版本不兼容会导致指标收集逻辑失效。检查任务运行状态
确保Flink任务已进入RUNNING状态,所有TaskManager都正常启动并注册到JobManager,初始化或重启中的任务无法提供完整指标。
内容的提问来源于stack exchange,提问作者crimson.blu
相关产品推荐
相关产品推荐

