Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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 String connectedNodes, final int connectedNodeCount, final int totalNodeCount) {
final String connectedNodes, 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, connectedNodes);
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();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8171,7 +8171,8 @@ protected Collection<AbstractMetricsRegistry> populateFlowMetrics(FlowMetricsRep
}
final boolean isClustered = clusterCoordinator != null;
final boolean isConnectedToCluster = isClustered() && clusterCoordinator.isConnected();
PrometheusMetricsUtil.createClusterMetrics(clusterMetricsRegistry, instanceId, isClustered, isConnectedToCluster, connectedNodesLabel, connectedNodeCount, totalNodeCount);
PrometheusMetricsUtil.createClusterMetrics(clusterMetricsRegistry, instanceId, isClustered, isConnectedToCluster, connectedNodesLabel, connectedNodeCount, totalNodeCount,
controllerFacade.isPrimary(), controllerFacade.isClusterCoordinator());
Collection<AbstractMetricsRegistry> metricsRegistries = Arrays.asList(
nifiMetricsRegistry,
jvmMetricsRegistry,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -276,13 +276,13 @@ public void testGetFlowMetricsPrometheusAsJson() throws IOException {
assertTrue(metrics.containsKey(ROOT_FIELD_NAME));

final List<Sample> registryList = metrics.get(ROOT_FIELD_NAME);
assertEquals(13, registryList.size());
assertEquals(15, registryList.size());

final Map<String, Long> 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
Expand Down Expand Up @@ -836,6 +836,8 @@ private static CollectorRegistry getClusterMetricsRegistry() {
clusterMetricsRegistry.setDataPoint(1, "IS_CONNECTED_TO_CLUSTER", "B1Id");
clusterMetricsRegistry.setDataPoint(2, "CONNECTED_NODE_COUNT", "B1Id", "2 / 3");
clusterMetricsRegistry.setDataPoint(3, "TOTAL_NODE_COUNT", "B1Id");
clusterMetricsRegistry.setDataPoint(1, "IS_PRIMARY_NODE", "B1Id");
clusterMetricsRegistry.setDataPoint(0, "IS_CLUSTER_COORDINATOR", "B1Id");

return clusterMetricsRegistry.getRegistry();
}
Expand Down
Loading