From 4887a109d10a3010fc8d1d6c93ce7310bee1fd98 Mon Sep 17 00:00:00 2001 From: Alexander Bij <1022013+abij@users.noreply.github.com> Date: Fri, 4 Sep 2026 16:20:13 +0200 Subject: [PATCH] NIFI-16298 Expose Primary Node and Cluster Coordinator role as Prometheus metrics --- .../prometheusutil/ClusterMetricsRegistry.java | 12 ++++++++++++ .../nifi/prometheusutil/PrometheusMetricsUtil.java | 5 ++++- .../apache/nifi/web/StandardNiFiServiceFacade.java | 3 ++- .../nifi/web/controller/ControllerFacade.java | 14 ++++++++++++++ .../org/apache/nifi/web/api/TestFlowResource.java | 6 ++++-- 5 files changed, 36 insertions(+), 4 deletions(-) diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/prometheusutil/ClusterMetricsRegistry.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/prometheusutil/ClusterMetricsRegistry.java index 462b2328812b..787b6e22bc5c 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/prometheusutil/ClusterMetricsRegistry.java +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/prometheusutil/ClusterMetricsRegistry.java @@ -48,5 +48,17 @@ public ClusterMetricsRegistry() { .help("The total number of nodes in this cluster") .labelNames("instance") .register(registry)); + + nameToGaugeMap.put("IS_PRIMARY_NODE", Gauge.build() + .name("cluster_is_primary_node") + .help("Whether this NiFi instance is the Primary Node. Values are 0 or 1") + .labelNames("instance") + .register(registry)); + + nameToGaugeMap.put("IS_CLUSTER_COORDINATOR", Gauge.build() + .name("cluster_is_cluster_coordinator") + .help("Whether this NiFi instance is the Cluster Coordinator. Values are 0 or 1") + .labelNames("instance") + .register(registry)); } } diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/prometheusutil/PrometheusMetricsUtil.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/prometheusutil/PrometheusMetricsUtil.java index 60c7d4d223af..1d799a0ab69f 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/prometheusutil/PrometheusMetricsUtil.java +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/prometheusutil/PrometheusMetricsUtil.java @@ -500,12 +500,15 @@ public static void createVersionInfoMetrics(final VersionInfoRegistry versionInf } public static CollectorRegistry createClusterMetrics(final ClusterMetricsRegistry clusterMetricsRegistry, final String instId, final boolean isClustered, final boolean isConnectedToCluster, - final int connectedNodeCount, final int totalNodeCount) { + final int connectedNodeCount, final int totalNodeCount, + final boolean isPrimaryNode, final boolean isClusterCoordinator) { final String instanceId = StringUtils.isEmpty(instId) ? DEFAULT_LABEL_STRING : instId; clusterMetricsRegistry.setDataPoint(isClustered ? 1 : 0, "IS_CLUSTERED", instanceId); clusterMetricsRegistry.setDataPoint(isConnectedToCluster ? 1 : 0, "IS_CONNECTED_TO_CLUSTER", instanceId); clusterMetricsRegistry.setDataPoint(connectedNodeCount, "CONNECTED_NODE_COUNT", instanceId); clusterMetricsRegistry.setDataPoint(totalNodeCount, "TOTAL_NODE_COUNT", instanceId); + clusterMetricsRegistry.setDataPoint(isPrimaryNode ? 1 : 0, "IS_PRIMARY_NODE", instanceId); + clusterMetricsRegistry.setDataPoint(isClusterCoordinator ? 1 : 0, "IS_CLUSTER_COORDINATOR", instanceId); return clusterMetricsRegistry.getRegistry(); } diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/StandardNiFiServiceFacade.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/StandardNiFiServiceFacade.java index 786b008514a2..00853d01fa40 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/StandardNiFiServiceFacade.java +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/StandardNiFiServiceFacade.java @@ -8166,7 +8166,8 @@ protected Collection populateFlowMetrics(FlowMetricsRep } final boolean isClustered = clusterCoordinator != null; final boolean isConnectedToCluster = isClustered() && clusterCoordinator.isConnected(); - PrometheusMetricsUtil.createClusterMetrics(clusterMetricsRegistry, instanceId, isClustered, isConnectedToCluster, connectedNodeCount, totalNodeCount); + PrometheusMetricsUtil.createClusterMetrics(clusterMetricsRegistry, instanceId, isClustered, isConnectedToCluster, connectedNodeCount, totalNodeCount, + controllerFacade.isPrimary(), controllerFacade.isClusterCoordinator()); Collection metricsRegistries = Arrays.asList( nifiMetricsRegistry, jvmMetricsRegistry, diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/controller/ControllerFacade.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/controller/ControllerFacade.java index e9a2436aa615..d33b4156dfc3 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/controller/ControllerFacade.java +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/main/java/org/apache/nifi/web/controller/ControllerFacade.java @@ -494,6 +494,20 @@ public boolean isClustered() { return flowController.isClustered(); } + /** + * @return true if this node is the Primary Node + */ + public boolean isPrimary() { + return flowController.isPrimary(); + } + + /** + * @return true if this node is the Cluster Coordinator + */ + public boolean isClusterCoordinator() { + return flowController.isClusterCoordinator(); + } + /** * Gets the name of this controller. * diff --git a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/TestFlowResource.java b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/TestFlowResource.java index 8356406c31ab..5a8c2187fe3a 100644 --- a/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/TestFlowResource.java +++ b/nifi-framework-bundle/nifi-framework/nifi-web/nifi-web-api/src/test/java/org/apache/nifi/web/api/TestFlowResource.java @@ -276,13 +276,13 @@ public void testGetFlowMetricsPrometheusAsJson() throws IOException { assertTrue(metrics.containsKey(ROOT_FIELD_NAME)); final List registryList = metrics.get(ROOT_FIELD_NAME); - assertEquals(13, registryList.size()); + assertEquals(15, registryList.size()); final Map result = getResult(registryList); assertEquals(3L, result.get(SAMPLE_NAME_JVM)); assertEquals(4L, result.get(SAMPLE_LABEL_VALUES_PROCESS_GROUP)); assertEquals(2L, result.get(SAMPLE_LABEL_VALUES_ROOT_PROCESS_GROUP)); - assertEquals(4L, result.get(CLUSTER_LABEL_KEY)); + assertEquals(6L, result.get(CLUSTER_LABEL_KEY)); } @Test @@ -836,6 +836,8 @@ private static CollectorRegistry getClusterMetricsRegistry() { clusterMetricsRegistry.setDataPoint(1, "IS_CONNECTED_TO_CLUSTER", "B1Id"); clusterMetricsRegistry.setDataPoint(2, "CONNECTED_NODE_COUNT", "B1Id"); clusterMetricsRegistry.setDataPoint(3, "TOTAL_NODE_COUNT", "B1Id"); + clusterMetricsRegistry.setDataPoint(1, "IS_PRIMARY_NODE", "B1Id"); + clusterMetricsRegistry.setDataPoint(0, "IS_CLUSTER_COORDINATOR", "B1Id"); return clusterMetricsRegistry.getRegistry(); }