diff --git a/modules/core/src/main/java/org/apache/ignite/internal/ClusterMetricsSnapshot.java b/modules/core/src/main/java/org/apache/ignite/internal/ClusterMetricsSnapshot.java
index 4d2795abffa01..7fb6ee5dd63ca 100644
--- a/modules/core/src/main/java/org/apache/ignite/internal/ClusterMetricsSnapshot.java
+++ b/modules/core/src/main/java/org/apache/ignite/internal/ClusterMetricsSnapshot.java
@@ -20,24 +20,24 @@
import java.nio.ByteBuffer;
import java.util.Collection;
import java.util.Map;
-import org.apache.ignite.cluster.ClusterGroup;
import org.apache.ignite.cluster.ClusterMetrics;
import org.apache.ignite.cluster.ClusterNode;
-import org.apache.ignite.internal.processors.cluster.NodeMetricsMessage;
import org.apache.ignite.internal.util.typedef.F;
import org.apache.ignite.internal.util.typedef.internal.S;
import org.apache.ignite.internal.util.typedef.internal.U;
+import org.apache.ignite.plugin.extensions.communication.Message;
import org.jetbrains.annotations.Nullable;
+import static java.lang.Math.max;
import static java.lang.Math.min;
/**
- * Implementation for {@link ClusterMetrics} interface.
+ * Implementation for {@link ClusterMetrics} interface which is also {@link Message}.
*
* Note that whenever adding or removing metric parameters, care
* must be taken to update serialize/deserialize logic as well.
*/
-public class ClusterMetricsSnapshot implements ClusterMetrics {
+public class ClusterMetricsSnapshot implements ClusterMetrics, Message {
/** Size of serialized node metrics. */
public static final int METRICS_SIZE =
4/*max active jobs*/ +
@@ -96,324 +96,1275 @@ public class ClusterMetricsSnapshot implements ClusterMetrics {
8/*current PME time*/;
/** */
- private NodeMetricsMessage m;
+ @Order(0)
+ public int maxActiveJobs = -1;
- /**
- * Creates empty snapshot.
- */
+ /** */
+ @Order(1)
+ public int curActiveJobs = -1;
+
+ /** */
+ @Order(2)
+ public float avgActiveJobs = -1;
+
+ /** */
+ @Order(3)
+ public int maxWaitingJobs = -1;
+
+ /** */
+ @Order(4)
+ public int curWaitingJobs = -1;
+
+ /** */
+ @Order(5)
+ public float avgWaitingJobs = -1;
+
+ /** */
+ @Order(6)
+ public int maxRejectedJobs = -1;
+
+ /** */
+ @Order(7)
+ public int curRejectedJobs = -1;
+
+ /** */
+ @Order(8)
+ public float avgRejectedJobs = -1;
+
+ /** */
+ @Order(9)
+ public int maxCancelledJobs = -1;
+
+ /** */
+ @Order(10)
+ public int curCancelledJobs = -1;
+
+ /** */
+ @Order(11)
+ public float avgCancelledJobs = -1;
+
+ /** */
+ @Order(12)
+ public int totalRejectedJobs = -1;
+
+ /** */
+ @Order(13)
+ public int totalCancelledJobs = -1;
+
+ /** */
+ @Order(14)
+ public int totalExecutedJobs = -1;
+
+ /** */
+ @Order(15)
+ public long maxJobWaitTime = -1;
+
+ /** */
+ @Order(16)
+ public long curJobWaitTime = Long.MAX_VALUE;
+
+ /** */
+ @Order(17)
+ public double avgJobWaitTime = -1;
+
+ /** */
+ @Order(18)
+ public long maxJobExecTime = -1;
+
+ /** */
+ @Order(19)
+ public long curJobExecTime = -1;
+
+ /** */
+ @Order(20)
+ public double avgJobExecTime = -1;
+
+ /** */
+ @Order(21)
+ public int totalExecTasks = -1;
+
+ /** */
+ @Order(22)
+ public long totalIdleTime = -1;
+
+ /** */
+ @Order(23)
+ public long curIdleTime = -1;
+
+ /** */
+ @Order(24)
+ public int totalCpus = -1;
+
+ /** */
+ @Order(25)
+ public double curCpuLoad = -1;
+
+ /** */
+ @Order(26)
+ public double avgCpuLoad = -1;
+
+ /** */
+ @Order(27)
+ public double curGcCpuLoad = -1;
+
+ /** */
+ @Order(28)
+ public long heapInit = -1;
+
+ /** */
+ @Order(29)
+ public long heapUsed = -1;
+
+ /** */
+ @Order(30)
+ public long heapCommitted = -1;
+
+ /** */
+ @Order(31)
+ public long heapMax = -1;
+
+ /** */
+ @Order(32)
+ public long heapTotal = -1;
+
+ /** */
+ @Order(33)
+ public long nonHeapInit = -1;
+
+ /** */
+ @Order(34)
+ public long nonHeapUsed = -1;
+
+ /** */
+ @Order(35)
+ public long nonHeapCommitted = -1;
+
+ /** */
+ @Order(36)
+ public long nonHeapMax = -1;
+
+ /** */
+ @Order(37)
+ public long nonHeapTotal = -1;
+
+ /** */
+ @Order(38)
+ public long upTime = -1;
+
+ /** */
+ @Order(39)
+ public long startTime = -1;
+
+ /** */
+ @Order(40)
+ public long nodeStartTime = -1;
+
+ /** */
+ @Order(41)
+ public int threadCnt = -1;
+
+ /** */
+ @Order(42)
+ public int peakThreadCnt = -1;
+
+ /** */
+ @Order(43)
+ public long startedThreadCnt = -1;
+
+ /** */
+ @Order(44)
+ public int daemonThreadCnt = -1;
+
+ /** */
+ @Order(45)
+ public long lastDataVer = -1;
+
+ /** */
+ @Order(46)
+ public int sentMsgsCnt = -1;
+
+ /** */
+ @Order(47)
+ public long sentBytesCnt = -1;
+
+ /** */
+ @Order(48)
+ public int rcvdMsgsCnt = -1;
+
+ /** */
+ @Order(49)
+ public long rcvdBytesCnt = -1;
+
+ /** */
+ @Order(50)
+ public int outMesQueueSize = -1;
+
+ /** */
+ @Order(51)
+ public int totalNodes = -1;
+
+ /** */
+ @Order(52)
+ public long totalJobsExecTime = -1;
+
+ /** */
+ @Order(53)
+ public long curPmeDuration = -1;
+
+ /** */
+ public long lastUpdateTime = -1;
+
+ /** Empty constructor for serialization purposes. */
public ClusterMetricsSnapshot() {
- m = new NodeMetricsMessage();
+ // Like in deserealize()
+ lastUpdateTime = System.currentTimeMillis();
}
/**
- * Creates snapshot based on the handled message.
+ * Create metrics for given nodes.
+ *
+ * @param nodes Nodes.
*/
- public ClusterMetricsSnapshot(NodeMetricsMessage m) {
- // As in #deserialize().
- m.lastUpdateTime = U.currentTimeMillis();
+ public ClusterMetricsSnapshot(Collection nodes) {
+ int size = nodes.size();
+
+ curJobWaitTime = Long.MAX_VALUE;
+ lastUpdateTime = 0;
+ maxActiveJobs = 0;
+ curActiveJobs = 0;
+ avgActiveJobs = 0;
+ maxWaitingJobs = 0;
+ curWaitingJobs = 0;
+ avgWaitingJobs = 0;
+ maxRejectedJobs = 0;
+ curRejectedJobs = 0;
+ avgRejectedJobs = 0;
+ maxCancelledJobs = 0;
+ curCancelledJobs = 0;
+ avgCancelledJobs = 0;
+ totalRejectedJobs = 0;
+ totalCancelledJobs = 0;
+ totalExecutedJobs = 0;
+ totalJobsExecTime = 0;
+ maxJobWaitTime = 0;
+ avgJobWaitTime = 0;
+ maxJobExecTime = 0;
+ curJobExecTime = 0;
+ avgJobExecTime = 0;
+ totalExecTasks = 0;
+ totalIdleTime = 0;
+ curIdleTime = 0;
+ totalCpus = 0;
+ curCpuLoad = 0;
+ avgCpuLoad = 0;
+ curGcCpuLoad = 0;
+ heapInit = 0;
+ heapUsed = 0;
+ heapCommitted = 0;
+ heapMax = 0;
+ nonHeapInit = 0;
+ nonHeapUsed = 0;
+ nonHeapCommitted = 0;
+ nonHeapMax = 0;
+ nonHeapTotal = 0;
+ upTime = 0;
+ startTime = 0;
+ nodeStartTime = 0;
+ threadCnt = 0;
+ peakThreadCnt = 0;
+ startedThreadCnt = 0;
+ daemonThreadCnt = 0;
+ lastDataVer = 0;
+ sentMsgsCnt = 0;
+ sentBytesCnt = 0;
+ rcvdMsgsCnt = 0;
+ rcvdBytesCnt = 0;
+ outMesQueueSize = 0;
+ heapTotal = 0;
+ totalNodes = nodes.size();
+ curPmeDuration = 0;
+
+ for (ClusterNode node : nodes) {
+ ClusterMetrics m = node.metrics();
+
+ lastUpdateTime = max(lastUpdateTime, node.metrics().getLastUpdateTime());
+
+ curActiveJobs += m.getCurrentActiveJobs();
+ maxActiveJobs = max(maxActiveJobs, m.getCurrentActiveJobs());
+ avgActiveJobs += m.getCurrentActiveJobs();
+ totalExecutedJobs += m.getTotalExecutedJobs();
+ totalJobsExecTime += m.getTotalJobsExecutionTime();
+
+ totalExecTasks += m.getTotalExecutedTasks();
+
+ totalCancelledJobs += m.getTotalCancelledJobs();
+ curCancelledJobs += m.getCurrentCancelledJobs();
+ maxCancelledJobs = max(maxCancelledJobs, m.getCurrentCancelledJobs());
+ avgCancelledJobs += m.getCurrentCancelledJobs();
+
+ totalRejectedJobs += m.getTotalRejectedJobs();
+ curRejectedJobs += m.getCurrentRejectedJobs();
+ maxRejectedJobs = max(maxRejectedJobs, m.getCurrentRejectedJobs());
+ avgRejectedJobs += m.getCurrentRejectedJobs();
+
+ curWaitingJobs += m.getCurrentWaitingJobs();
+ maxWaitingJobs = max(maxWaitingJobs, m.getCurrentWaitingJobs());
+ avgWaitingJobs += m.getCurrentWaitingJobs();
+
+ maxJobExecTime = max(maxJobExecTime, m.getMaximumJobExecuteTime());
+ avgJobExecTime += m.getAverageJobExecuteTime();
+ curJobExecTime += m.getCurrentJobExecuteTime();
+
+ curJobWaitTime = min(curJobWaitTime, m.getCurrentJobWaitTime());
+ maxJobWaitTime = max(maxJobWaitTime, m.getCurrentJobWaitTime());
+ avgJobWaitTime += m.getAverageJobWaitTime();
+
+ daemonThreadCnt += m.getCurrentDaemonThreadCount();
+
+ peakThreadCnt = max(peakThreadCnt, m.getCurrentThreadCount());
+ threadCnt += m.getCurrentThreadCount();
+ startedThreadCnt += m.getTotalStartedThreadCount();
+
+ curIdleTime += m.getCurrentIdleTime();
+ totalIdleTime += m.getTotalIdleTime();
+
+ heapCommitted += m.getHeapMemoryCommitted();
+
+ heapUsed += m.getHeapMemoryUsed();
+
+ heapMax = max(heapMax, m.getHeapMemoryMaximum());
+
+ heapTotal += m.getHeapMemoryTotal();
+
+ heapInit += m.getHeapMemoryInitialized();
+
+ nonHeapCommitted += m.getNonHeapMemoryCommitted();
+
+ nonHeapUsed += m.getNonHeapMemoryUsed();
+
+ nonHeapMax = max(nonHeapMax, m.getNonHeapMemoryMaximum());
+
+ nonHeapTotal += m.getNonHeapMemoryTotal();
+
+ nonHeapInit += m.getNonHeapMemoryInitialized();
+
+ upTime = max(upTime, m.getUpTime());
+
+ lastDataVer = max(lastDataVer, m.getLastDataVersion());
+
+ sentMsgsCnt += m.getSentMessagesCount();
+ sentBytesCnt += m.getSentBytesCount();
+ rcvdMsgsCnt += m.getReceivedMessagesCount();
+ rcvdBytesCnt += m.getReceivedBytesCount();
+ outMesQueueSize += m.getOutboundMessagesQueueSize();
+
+ avgCpuLoad += m.getCurrentCpuLoad();
+
+ curPmeDuration = max(curPmeDuration, m.getCurrentPmeDuration());
+ }
+
+ curJobExecTime /= size;
+
+ avgActiveJobs /= size;
+ avgCancelledJobs /= size;
+ avgRejectedJobs /= size;
+ avgWaitingJobs /= size;
+ avgJobExecTime /= size;
+ avgJobWaitTime /= size;
+ avgCpuLoad /= size;
+
+ if (!F.isEmpty(nodes)) {
+ ClusterMetrics oldestNodeMetrics = oldest(nodes).metrics();
+
+ nodeStartTime = oldestNodeMetrics.getNodeStartTime();
+ startTime = oldestNodeMetrics.getStartTime();
+ }
- this.m = m;
+ Map> neighborhood = U.neighborhood(nodes);
+
+ curGcCpuLoad = currentGcCpuLoad(neighborhood);
+ curCpuLoad = currentCpuLoad(neighborhood);
+ totalCpus = cpuCnt(neighborhood);
}
- /**
- * Create metrics for given cluster group.
- *
- * @param p Projection to get metrics for.
- */
- public ClusterMetricsSnapshot(ClusterGroup p) {
- assert p != null;
+ /** */
+ public ClusterMetricsSnapshot(ClusterMetrics metrics) {
+ maxActiveJobs = metrics.getMaximumActiveJobs();
+ curActiveJobs = metrics.getCurrentActiveJobs();
+ avgActiveJobs = metrics.getAverageActiveJobs();
+
+ maxWaitingJobs = metrics.getMaximumWaitingJobs();
+ curWaitingJobs = metrics.getCurrentWaitingJobs();
+ avgWaitingJobs = metrics.getAverageWaitingJobs();
+
+ maxRejectedJobs = metrics.getMaximumRejectedJobs();
+ curRejectedJobs = metrics.getCurrentRejectedJobs();
+ avgRejectedJobs = metrics.getAverageRejectedJobs();
+
+ maxCancelledJobs = metrics.getMaximumCancelledJobs();
+ curCancelledJobs = metrics.getCurrentCancelledJobs();
+ avgCancelledJobs = metrics.getAverageCancelledJobs();
+
+ totalRejectedJobs = metrics.getTotalRejectedJobs();
+ totalCancelledJobs = metrics.getTotalCancelledJobs();
+ totalExecutedJobs = metrics.getTotalExecutedJobs();
+
+ maxJobWaitTime = metrics.getMaximumJobWaitTime();
+ curJobWaitTime = metrics.getCurrentJobWaitTime();
+ avgJobWaitTime = metrics.getAverageJobWaitTime();
+
+ maxJobExecTime = metrics.getMaximumJobExecuteTime();
+ curJobExecTime = metrics.getCurrentJobExecuteTime();
+ avgJobExecTime = metrics.getAverageJobExecuteTime();
+
+ totalJobsExecTime = metrics.getTotalJobsExecutionTime();
+ totalExecTasks = metrics.getTotalExecutedTasks();
+
+ curIdleTime = metrics.getCurrentIdleTime();
+ totalIdleTime = metrics.getTotalIdleTime();
+
+ totalCpus = metrics.getTotalCpus();
+ curCpuLoad = metrics.getCurrentCpuLoad();
+ avgCpuLoad = metrics.getAverageCpuLoad();
+ curGcCpuLoad = metrics.getCurrentGcCpuLoad();
+
+ heapInit = metrics.getHeapMemoryInitialized();
+ heapUsed = metrics.getHeapMemoryUsed();
+ heapCommitted = metrics.getHeapMemoryCommitted();
+ heapMax = metrics.getHeapMemoryMaximum();
+ heapTotal = metrics.getHeapMemoryTotal();
- m = new NodeMetricsMessage(p.nodes());
+ nonHeapInit = metrics.getNonHeapMemoryInitialized();
+ nonHeapUsed = metrics.getNonHeapMemoryUsed();
+ nonHeapCommitted = metrics.getNonHeapMemoryCommitted();
+ nonHeapMax = metrics.getNonHeapMemoryMaximum();
+ nonHeapTotal = metrics.getNonHeapMemoryTotal();
+
+ startTime = metrics.getStartTime();
+ nodeStartTime = metrics.getNodeStartTime();
+ upTime = metrics.getUpTime();
+
+ lastDataVer = metrics.getLastDataVersion();
+
+ curPmeDuration = metrics.getCurrentPmeDuration();
+
+ totalNodes = metrics.getTotalNodes();
+
+ threadCnt = metrics.getCurrentThreadCount();
+ peakThreadCnt = metrics.getMaximumThreadCount();
+ startedThreadCnt = metrics.getTotalStartedThreadCount();
+ daemonThreadCnt = metrics.getCurrentDaemonThreadCount();
+
+ sentMsgsCnt = metrics.getSentMessagesCount();
+ rcvdMsgsCnt = metrics.getReceivedMessagesCount();
+ outMesQueueSize = metrics.getOutboundMessagesQueueSize();
+
+ sentBytesCnt = metrics.getSentBytesCount();
+ rcvdBytesCnt = metrics.getReceivedBytesCount();
+
+ lastUpdateTime = metrics.getLastUpdateTime();
+ }
+
+ /** */
+ public static ClusterMetricsSnapshot of(ClusterMetrics metrics) {
+ return metrics instanceof ClusterMetricsSnapshot ? (ClusterMetricsSnapshot)metrics : new ClusterMetricsSnapshot(metrics);
}
/** {@inheritDoc} */
@Override public long getHeapMemoryTotal() {
- return m.heapMemoryTotal();
+ return heapTotal;
+ }
+
+ /**
+ * Sets total heap size.
+ *
+ * @param heapTotal Total heap.
+ */
+ public void heapMemoryTotal(long heapTotal) {
+ this.heapTotal = heapTotal;
+ }
+
+ /**
+ * Sets non-heap total heap size.
+ *
+ * @param nonHeapTotal Total heap.
+ */
+ public void nonHeapMemoryTotal(long nonHeapTotal) {
+ this.nonHeapTotal = nonHeapTotal;
}
/** {@inheritDoc} */
@Override public long getLastUpdateTime() {
- return m.lastUpdateTime();
+ return lastUpdateTime;
+ }
+
+ /**
+ * Sets last update time.
+ *
+ * @param lastUpdateTime Last update time.
+ */
+ public void lastUpdateTime(long lastUpdateTime) {
+ this.lastUpdateTime = lastUpdateTime;
}
/** {@inheritDoc} */
@Override public int getMaximumActiveJobs() {
- return m.maximumActiveJobs();
+ return maxActiveJobs;
+ }
+
+ /**
+ * Sets max active jobs.
+ *
+ * @param maxActiveJobs Max active jobs.
+ */
+ public void maximumActiveJobs(int maxActiveJobs) {
+ this.maxActiveJobs = maxActiveJobs;
}
/** {@inheritDoc} */
@Override public int getCurrentActiveJobs() {
- return m.currentActiveJobs();
+ return curActiveJobs;
+ }
+
+ /**
+ * Sets current active jobs.
+ *
+ * @param curActiveJobs Current active jobs.
+ */
+ public void currentActiveJobs(int curActiveJobs) {
+ this.curActiveJobs = curActiveJobs;
}
/** {@inheritDoc} */
@Override public float getAverageActiveJobs() {
- return m.averageActiveJobs();
+ return curActiveJobs;
+ }
+
+ /**
+ * Sets average active jobs.
+ *
+ * @param avgActiveJobs Average active jobs.
+ */
+ public void averageActiveJobs(float avgActiveJobs) {
+ this.avgActiveJobs = avgActiveJobs;
}
/** {@inheritDoc} */
@Override public int getMaximumWaitingJobs() {
- return m.maximumWaitingJobs();
+ return maxWaitingJobs;
+ }
+
+ /**
+ * Sets maximum waiting jobs.
+ *
+ * @param maxWaitingJobs Maximum waiting jobs.
+ */
+ public void maximumWaitingJobs(int maxWaitingJobs) {
+ this.maxWaitingJobs = maxWaitingJobs;
}
/** {@inheritDoc} */
@Override public int getCurrentWaitingJobs() {
- return m.currentWaitingJobs();
+ return curWaitingJobs;
+ }
+
+ /**
+ * Sets current waiting jobs.
+ *
+ * @param curWaitingJobs Current waiting jobs.
+ */
+ public void currentWaitingJobs(int curWaitingJobs) {
+ this.curWaitingJobs = curWaitingJobs;
}
/** {@inheritDoc} */
@Override public float getAverageWaitingJobs() {
- return m.averageWaitingJobs();
+ return avgWaitingJobs;
+ }
+
+ /**
+ * Sets average waiting jobs.
+ *
+ * @param avgWaitingJobs Average waiting jobs.
+ */
+ public void averageWaitingJobs(float avgWaitingJobs) {
+ this.avgWaitingJobs = avgWaitingJobs;
}
/** {@inheritDoc} */
@Override public int getMaximumRejectedJobs() {
- return m.maximumRejectedJobs();
+ return maxRejectedJobs;
+ }
+
+ /**
+ * @param maxRejectedJobs Maximum number of jobs rejected during a single collision resolution event.
+ */
+ public void maximumRejectedJobs(int maxRejectedJobs) {
+ this.maxRejectedJobs = maxRejectedJobs;
}
/** {@inheritDoc} */
@Override public int getCurrentRejectedJobs() {
- return m.currentRejectedJobs();
+ return curRejectedJobs;
+ }
+
+ /**
+ * @param curRejectedJobs Number of jobs rejected during most recent collision resolution.
+ */
+ public void currentRejectedJobs(int curRejectedJobs) {
+ this.curRejectedJobs = curRejectedJobs;
}
/** {@inheritDoc} */
@Override public float getAverageRejectedJobs() {
- return m.averageRejectedJobs();
+ return avgRejectedJobs;
+ }
+
+ /**
+ * @param avgRejectedJobs Average number of jobs this node rejects.
+ */
+ public void averageRejectedJobs(float avgRejectedJobs) {
+ this.avgRejectedJobs = avgRejectedJobs;
}
/** {@inheritDoc} */
@Override public int getTotalRejectedJobs() {
- return m.totalRejectedJobs();
+ return totalRejectedJobs;
+ }
+
+ /**
+ * @param totalRejectedJobs Total number of jobs this node ever rejected.
+ */
+ public void totalRejectedJobs(int totalRejectedJobs) {
+ this.totalRejectedJobs = totalRejectedJobs;
}
/** {@inheritDoc} */
@Override public int getMaximumCancelledJobs() {
- return m.maximumCancelledJobs();
+ return maxCancelledJobs;
+ }
+
+ /**
+ * Sets maximum cancelled jobs.
+ *
+ * @param maxCancelledJobs Maximum cancelled jobs.
+ */
+ public void maximumCancelledJobs(int maxCancelledJobs) {
+ this.maxCancelledJobs = maxCancelledJobs;
}
/** {@inheritDoc} */
@Override public int getCurrentCancelledJobs() {
- return m.currentCancelledJobs();
+ return curCancelledJobs;
+ }
+
+ /**
+ * Sets current cancelled jobs.
+ *
+ * @param curCancelledJobs Current cancelled jobs.
+ */
+ public void currentCancelledJobs(int curCancelledJobs) {
+ this.curCancelledJobs = curCancelledJobs;
}
/** {@inheritDoc} */
@Override public float getAverageCancelledJobs() {
- return m.averageCancelledJobs();
+ return avgCancelledJobs;
+ }
+
+ /**
+ * Sets average cancelled jobs.
+ *
+ * @param avgCancelledJobs Average cancelled jobs.
+ */
+ public void averageCancelledJobs(float avgCancelledJobs) {
+ this.avgCancelledJobs = avgCancelledJobs;
}
/** {@inheritDoc} */
@Override public int getTotalExecutedJobs() {
- return m.totalExecutedJobs();
+ return totalExecutedJobs;
+ }
+
+ /**
+ * Sets total active jobs.
+ *
+ * @param totalExecutedJobs Total active jobs.
+ */
+ public void totalExecutedJobs(int totalExecutedJobs) {
+ this.totalExecutedJobs = totalExecutedJobs;
}
/** {@inheritDoc} */
@Override public long getTotalJobsExecutionTime() {
- return m.totalJobsExecutionTime();
+ return totalJobsExecTime;
+ }
+
+ /**
+ * Sets total jobs execution time.
+ *
+ * @param totalJobsExecTime Total jobs execution time.
+ */
+ public void totalJobsExecutionTime(long totalJobsExecTime) {
+ this.totalJobsExecTime = totalJobsExecTime;
}
/** {@inheritDoc} */
@Override public int getTotalCancelledJobs() {
- return m.totalCancelledJobs();
+ return totalCancelledJobs;
+ }
+
+ /**
+ * Sets total cancelled jobs.
+ *
+ * @param totalCancelledJobs Total cancelled jobs.
+ */
+ public void totalCancelledJobs(int totalCancelledJobs) {
+ this.totalCancelledJobs = totalCancelledJobs;
}
/** {@inheritDoc} */
@Override public long getMaximumJobWaitTime() {
- return m.maximumJobWaitTime();
+ return maxJobWaitTime;
+ }
+
+ /**
+ * Sets max job wait time.
+ *
+ * @param maxJobWaitTime Max job wait time.
+ */
+ public void maximumJobWaitTime(long maxJobWaitTime) {
+ this.maxJobWaitTime = maxJobWaitTime;
}
/** {@inheritDoc} */
@Override public long getCurrentJobWaitTime() {
- return m.currentJobWaitTime();
+ return curJobWaitTime;
+ }
+
+ /**
+ * Sets current job wait time.
+ *
+ * @param curJobWaitTime Current job wait time.
+ */
+ public void currentJobWaitTime(long curJobWaitTime) {
+ this.curJobWaitTime = curJobWaitTime;
}
/** {@inheritDoc} */
@Override public double getAverageJobWaitTime() {
- return m.averageJobWaitTime();
+ return avgJobWaitTime;
+ }
+
+ /**
+ * Sets average job wait time.
+ *
+ * @param avgJobWaitTime Average job wait time.
+ */
+ public void averageJobWaitTime(double avgJobWaitTime) {
+ this.avgJobWaitTime = avgJobWaitTime;
}
/** {@inheritDoc} */
@Override public long getMaximumJobExecuteTime() {
- return m.maximumJobExecuteTime();
+ return maxJobExecTime;
+ }
+
+ /**
+ * Sets maximum job execution time.
+ *
+ * @param maxJobExecTime Maximum job execution time.
+ */
+ public void maximumJobExecuteTime(long maxJobExecTime) {
+ this.maxJobExecTime = maxJobExecTime;
}
/** {@inheritDoc} */
@Override public long getCurrentJobExecuteTime() {
- return m.currentJobExecuteTime();
+ return curJobExecTime;
+ }
+
+ /**
+ * Sets current job execute time.
+ *
+ * @param curJobExecTime Current job execute time.
+ */
+ public void currentJobExecuteTime(long curJobExecTime) {
+ this.curJobExecTime = curJobExecTime;
}
/** {@inheritDoc} */
@Override public double getAverageJobExecuteTime() {
- return m.averageJobExecuteTime();
+ return avgJobExecTime;
+ }
+
+ /**
+ * Sets average job execution time.
+ *
+ * @param avgJobExecTime Average job execution time.
+ */
+ public void averageJobExecuteTime(double avgJobExecTime) {
+ this.avgJobExecTime = avgJobExecTime;
}
/** {@inheritDoc} */
@Override public int getTotalExecutedTasks() {
- return m.totalExecutedTasks();
+ return totalExecTasks;
}
- /** {@inheritDoc} */
- @Override public long getTotalBusyTime() {
- return getUpTime() - getTotalIdleTime();
+ /**
+ * Sets total executed tasks count.
+ *
+ * @param totalExecTasks total executed tasks count.
+ */
+ public void totalExecutedTasks(int totalExecTasks) {
+ this.totalExecTasks = totalExecTasks;
}
/** {@inheritDoc} */
@Override public long getTotalIdleTime() {
- return m.totalIdleTime();
+ return totalIdleTime;
}
- /** {@inheritDoc} */
- @Override public long getCurrentIdleTime() {
- return m.currentIdleTime();
+ /**
+ * Set total node idle time.
+ *
+ * @param totalIdleTime Total node idle time.
+ */
+ public void totalIdleTime(long totalIdleTime) {
+ this.totalIdleTime = totalIdleTime;
}
/** {@inheritDoc} */
- @Override public float getBusyTimePercentage() {
- return 1 - getIdleTimePercentage();
+ @Override public long getCurrentIdleTime() {
+ return curIdleTime;
}
- /** {@inheritDoc} */
- @Override public float getIdleTimePercentage() {
- return getTotalIdleTime() / (float)getUpTime();
+ /**
+ * Sets time elapsed since execution of last job.
+ *
+ * @param curIdleTime Time elapsed since execution of last job.
+ */
+ public void currentIdleTime(long curIdleTime) {
+ this.curIdleTime = curIdleTime;
}
/** {@inheritDoc} */
@Override public int getTotalCpus() {
- return m.totalCpus();
+ return totalCpus;
}
/** {@inheritDoc} */
@Override public double getCurrentCpuLoad() {
- return m.currentCpuLoad();
+ return curCpuLoad;
}
/** {@inheritDoc} */
@Override public double getAverageCpuLoad() {
- return m.averageCpuLoad();
+ return avgCpuLoad;
}
/** {@inheritDoc} */
@Override public double getCurrentGcCpuLoad() {
- return m.currentGcCpuLoad();
+ return curGcCpuLoad;
}
/** {@inheritDoc} */
@Override public long getHeapMemoryInitialized() {
- return m.heapMemoryInitialized();
+ return heapInit;
}
/** {@inheritDoc} */
@Override public long getHeapMemoryUsed() {
- return m.heapMemoryUsed();
+ return heapUsed;
}
/** {@inheritDoc} */
@Override public long getHeapMemoryCommitted() {
- return m.heapMemoryCommitted();
+ return heapCommitted;
}
/** {@inheritDoc} */
@Override public long getHeapMemoryMaximum() {
- return m.heapMemoryMaximum();
+ return heapMax;
}
/** {@inheritDoc} */
@Override public long getNonHeapMemoryInitialized() {
- return m.nonHeapMemoryInitialized();
+ return nonHeapInit;
}
/** {@inheritDoc} */
@Override public long getNonHeapMemoryUsed() {
- return m.nonHeapMemoryUsed();
+ return nonHeapUsed;
}
/** {@inheritDoc} */
@Override public long getNonHeapMemoryCommitted() {
- return m.nonHeapMemoryCommitted();
+ return nonHeapCommitted;
}
/** {@inheritDoc} */
@Override public long getNonHeapMemoryMaximum() {
- return m.nonHeapMemoryMaximum();
+ return nonHeapMax;
}
/** {@inheritDoc} */
@Override public long getNonHeapMemoryTotal() {
- return m.nonHeapMemoryTotal();
+ return nonHeapTotal;
}
/** {@inheritDoc} */
@Override public long getUpTime() {
- return m.upTime();
+ return upTime;
}
/** {@inheritDoc} */
@Override public long getStartTime() {
- return m.startTime();
+ return startTime;
}
/** {@inheritDoc} */
@Override public long getNodeStartTime() {
- return m.nodeStartTime();
+ return nodeStartTime;
}
/** {@inheritDoc} */
@Override public int getCurrentThreadCount() {
- return m.currentThreadCount();
+ return threadCnt;
}
/** {@inheritDoc} */
@Override public int getMaximumThreadCount() {
- return m.maximumThreadCount();
+ return peakThreadCnt;
}
/** {@inheritDoc} */
@Override public long getTotalStartedThreadCount() {
- return m.totalStartedThreadCount();
+ return startedThreadCnt;
}
/** {@inheritDoc} */
@Override public int getCurrentDaemonThreadCount() {
- return m.currentDaemonThreadCount();
+ return daemonThreadCnt;
}
/** {@inheritDoc} */
@Override public long getLastDataVersion() {
- return m.lastDataVersion();
+ return lastDataVer;
}
/** {@inheritDoc} */
@Override public int getSentMessagesCount() {
- return m.sentMessagesCount();
+ return sentMsgsCnt;
}
/** {@inheritDoc} */
@Override public long getSentBytesCount() {
- return m.sentBytesCount();
+ return sentBytesCnt;
}
/** {@inheritDoc} */
@Override public int getReceivedMessagesCount() {
- return m.receivedMessagesCount();
+ return rcvdMsgsCnt;
}
/** {@inheritDoc} */
@Override public long getReceivedBytesCount() {
- return m.receivedBytesCount();
+ return rcvdBytesCnt;
}
/** {@inheritDoc} */
@Override public int getOutboundMessagesQueueSize() {
- return m.outboundMessagesQueueSize();
+ return outMesQueueSize;
}
/** {@inheritDoc} */
@Override public int getTotalNodes() {
- return m.totalNodes();
+ return totalNodes;
}
/** {@inheritDoc} */
@Override public long getCurrentPmeDuration() {
- return m.currentPmeDuration();
+ return curPmeDuration;
+ }
+
+ /** {@inheritDoc} */
+ @Override public long getTotalBusyTime() {
+ return getUpTime() - getTotalIdleTime();
+ }
+
+ /** {@inheritDoc} */
+ @Override public float getBusyTimePercentage() {
+ return 1 - getIdleTimePercentage();
+ }
+
+ /** {@inheritDoc} */
+ @Override public float getIdleTimePercentage() {
+ return getTotalIdleTime() / (float)getUpTime();
+ }
+
+ /**
+ * Sets available processors.
+ *
+ * @param totalCpus Available processors.
+ */
+ public void totalCpus(int totalCpus) {
+ this.totalCpus = totalCpus;
+ }
+
+ /**
+ * Sets current CPU load.
+ *
+ * @param curCpuLoad Current CPU load.
+ */
+ public void currentCpuLoad(double curCpuLoad) {
+ this.curCpuLoad = curCpuLoad;
+ }
+
+ /**
+ * Sets CPU load average over the metrics history.
+ *
+ * @param avgCpuLoad CPU load average.
+ */
+ public void averageCpuLoad(double avgCpuLoad) {
+ this.avgCpuLoad = avgCpuLoad;
+ }
+
+ /**
+ * Sets current GC load.
+ *
+ * @param curGcCpuLoad Current GC load.
+ */
+ public void currentGcCpuLoad(double curGcCpuLoad) {
+ this.curGcCpuLoad = curGcCpuLoad;
+ }
+
+ /**
+ * Sets heap initial memory.
+ *
+ * @param heapInit Heap initial memory.
+ */
+ public void heapMemoryInitialized(long heapInit) {
+ this.heapInit = heapInit;
+ }
+
+ /**
+ * Sets used heap memory.
+ *
+ * @param heapUsed Used heap memory.
+ */
+ public void heapMemoryUsed(long heapUsed) {
+ this.heapUsed = heapUsed;
+ }
+
+ /**
+ * Sets committed heap memory.
+ *
+ * @param heapCommitted Committed heap memory.
+ */
+ public void heapMemoryCommitted(long heapCommitted) {
+ this.heapCommitted = heapCommitted;
+ }
+
+ /**
+ * Sets maximum possible heap memory.
+ *
+ * @param heapMax Maximum possible heap memory.
+ */
+ public void heapMemoryMaximum(long heapMax) {
+ this.heapMax = heapMax;
+ }
+
+ /**
+ * Sets initial non-heap memory.
+ *
+ * @param nonHeapInit Initial non-heap memory.
+ */
+ public void nonHeapMemoryInitialized(long nonHeapInit) {
+ this.nonHeapInit = nonHeapInit;
+ }
+
+ /**
+ * Sets used non-heap memory.
+ *
+ * @param nonHeapUsed Used non-heap memory.
+ */
+ public void nonHeapMemoryUsed(long nonHeapUsed) {
+ this.nonHeapUsed = nonHeapUsed;
+ }
+
+ /**
+ * Sets committed non-heap memory.
+ *
+ * @param nonHeapCommitted Committed non-heap memory.
+ */
+ public void nonHeapMemoryCommitted(long nonHeapCommitted) {
+ this.nonHeapCommitted = nonHeapCommitted;
+ }
+
+ /**
+ * Sets maximum possible non-heap memory.
+ *
+ * @param nonHeapMax Maximum possible non-heap memory.
+ */
+ public void nonHeapMemoryMaximum(long nonHeapMax) {
+ this.nonHeapMax = nonHeapMax;
+ }
+
+ /**
+ * Sets VM up time.
+ *
+ * @param upTime VM up time.
+ */
+ public void upTime(long upTime) {
+ this.upTime = upTime;
+ }
+
+ /**
+ * Sets VM start time.
+ *
+ * @param startTime VM start time.
+ */
+ public void startTime(long startTime) {
+ this.startTime = startTime;
+ }
+
+ /**
+ * Sets node start time.
+ *
+ * @param nodeStartTime node start time.
+ */
+ public void nodeStartTime(long nodeStartTime) {
+ this.nodeStartTime = nodeStartTime;
+ }
+
+ /**
+ * Sets thread count.
+ *
+ * @param threadCnt Thread count.
+ */
+ public void currentThreadCount(int threadCnt) {
+ this.threadCnt = threadCnt;
+ }
+
+ /**
+ * Sets peak thread count.
+ *
+ * @param peakThreadCnt Peak thread count.
+ */
+ public void maximumThreadCount(int peakThreadCnt) {
+ this.peakThreadCnt = peakThreadCnt;
+ }
+
+ /**
+ * Sets started thread count.
+ *
+ * @param startedThreadCnt Started thread count.
+ */
+ public void totalStartedThreadCount(long startedThreadCnt) {
+ this.startedThreadCnt = startedThreadCnt;
+ }
+
+ /**
+ * Sets daemon thread count.
+ *
+ * @param daemonThreadCnt Daemon thread count.
+ */
+ public void currentDaemonThreadCount(int daemonThreadCnt) {
+ this.daemonThreadCnt = daemonThreadCnt;
+ }
+
+ /**
+ * Sets last data version.
+ *
+ * @param lastDataVer Last data version.
+ */
+ public void lastDataVersion(long lastDataVer) {
+ this.lastDataVer = lastDataVer;
+ }
+
+ /**
+ * Sets sent messages count.
+ *
+ * @param sentMsgsCnt Sent messages count.
+ */
+ public void sentMessagesCount(int sentMsgsCnt) {
+ this.sentMsgsCnt = sentMsgsCnt;
+ }
+
+ /**
+ * Sets sent bytes count.
+ *
+ * @param sentBytesCnt Sent bytes count.
+ */
+ public void sentBytesCount(long sentBytesCnt) {
+ this.sentBytesCnt = sentBytesCnt;
+ }
+
+ /**
+ * Sets received messages count.
+ *
+ * @param rcvdMsgsCnt Received messages count.
+ */
+ public void receivedMessagesCount(int rcvdMsgsCnt) {
+ this.rcvdMsgsCnt = rcvdMsgsCnt;
+ }
+
+ /**
+ * Sets received bytes count.
+ *
+ * @param rcvdBytesCnt Received bytes count.
+ */
+ public void receivedBytesCount(long rcvdBytesCnt) {
+ this.rcvdBytesCnt = rcvdBytesCnt;
+ }
+
+ /**
+ * Sets outbound messages queue size.
+ *
+ * @param outMesQueueSize Outbound messages queue size.
+ */
+ public void outboundMessagesQueueSize(int outMesQueueSize) {
+ this.outMesQueueSize = outMesQueueSize;
+ }
+
+ /**
+ * Sets total number of nodes.
+ *
+ * @param totalNodes Total number of nodes.
+ */
+ public void totalNodes(int totalNodes) {
+ this.totalNodes = totalNodes;
+ }
+
+ /**
+ * Sets execution duration for current partition map exchange.
+ *
+ * @param curPmeDuration Execution duration for current partition map exchange.
+ */
+ public void currentPmeDuration(long curPmeDuration) {
+ this.curPmeDuration = curPmeDuration;
+ }
+
+ /**
+ * Gets the oldest node in given collection.
+ *
+ * @param nodes Nodes.
+ * @return Oldest node or {@code null} if collection is empty.
+ */
+ @Nullable private static ClusterNode oldest(Collection nodes) {
+ long min = Long.MAX_VALUE;
+
+ ClusterNode oldest = null;
+
+ for (ClusterNode n : nodes)
+ if (n.order() < min) {
+ min = n.order();
+ oldest = n;
+ }
+
+ return oldest;
}
/**
@@ -438,56 +1389,36 @@ private static int cpuCnt(Map> neighborhood) {
* @param neighborhood Cluster neighborhood.
* @return CPU load.
*/
- private static int cpus(Map> neighborhood) {
- int cpus = 0;
+ private static double currentCpuLoad(Map> neighborhood) {
+ double curCpuLoad = 0.0;
for (Collection nodes : neighborhood.values()) {
ClusterNode first = F.first(nodes);
// Projection can be empty if all nodes in it failed.
if (first != null)
- cpus += first.metrics().getCurrentCpuLoad();
+ curCpuLoad += first.metrics().getCurrentCpuLoad();
}
- return cpus;
+ return curCpuLoad;
}
/**
* @param neighborhood Cluster neighborhood.
* @return GC CPU load.
*/
- private static int gcCpus(Map> neighborhood) {
- int cpus = 0;
+ private static double currentGcCpuLoad(Map> neighborhood) {
+ double curGcCpuLoad = 0;
for (Collection nodes : neighborhood.values()) {
ClusterNode first = F.first(nodes);
// Projection can be empty if all nodes in it failed.
if (first != null)
- cpus += first.metrics().getCurrentGcCpuLoad();
+ curGcCpuLoad += first.metrics().getCurrentGcCpuLoad();
}
- return cpus;
- }
-
- /**
- * Gets the oldest node in given collection.
- *
- * @param nodes Nodes.
- * @return Oldest node or {@code null} if collection is empty.
- */
- @Nullable private static ClusterNode oldest(Collection nodes) {
- long min = Long.MAX_VALUE;
-
- ClusterNode oldest = null;
-
- for (ClusterNode n : nodes)
- if (n.order() < min) {
- min = n.order();
- oldest = n;
- }
-
- return oldest;
+ return curGcCpuLoad;
}
/**
@@ -583,8 +1514,8 @@ public static int serialize(byte[] data, int off, ClusterMetrics metrics) {
* @param off Offset into byte array.
* @return Deserialized node metrics.
*/
- public static ClusterMetrics deserialize(byte[] data, int off) {
- NodeMetricsMessage msg = new NodeMetricsMessage();
+ public static ClusterMetricsSnapshot deserialize(byte[] data, int off) {
+ ClusterMetricsSnapshot msg = new ClusterMetricsSnapshot();
int bufSize = min(METRICS_SIZE, data.length - off);
@@ -656,7 +1587,7 @@ public static ClusterMetrics deserialize(byte[] data, int off) {
else
msg.currentPmeDuration(0);
- return new ClusterMetricsSnapshot(msg);
+ return msg;
}
/** {@inheritDoc} */
diff --git a/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java b/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java
index a460ed789d24b..16412f3f5a67d 100644
--- a/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java
+++ b/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java
@@ -199,7 +199,6 @@
import org.apache.ignite.internal.processors.cluster.ClusterUpdateNotifierDataBagItem;
import org.apache.ignite.internal.processors.cluster.DiscoveryDataClusterState;
import org.apache.ignite.internal.processors.cluster.NodeFullMetricsMessage;
-import org.apache.ignite.internal.processors.cluster.NodeMetricsMessage;
import org.apache.ignite.internal.processors.continuous.ContinuousRoutineStartResultMessage;
import org.apache.ignite.internal.processors.continuous.GridContinuousMessage;
import org.apache.ignite.internal.processors.continuous.StartRequestData;
@@ -666,7 +665,7 @@ public CoreMessagesProvider(Marshaller dfltMarsh, Marshaller schemaAwareMarsh) {
// [11900 - 12000]: Metrics, monitoring messages.
msgIdx = 11900;
register(CacheMetricsMessage.class);
- register(NodeMetricsMessage.class);
+ register(ClusterMetricsSnapshot.class);
register(NodeFullMetricsMessage.class);
register(ClusterMetricsUpdateMessage.class);
register(TcpDiscoveryClientNodesMetricsMessage.class);
diff --git a/modules/core/src/main/java/org/apache/ignite/internal/cluster/ClusterGroupAdapter.java b/modules/core/src/main/java/org/apache/ignite/internal/cluster/ClusterGroupAdapter.java
index 33d833bfe8953..501a2df43590a 100644
--- a/modules/core/src/main/java/org/apache/ignite/internal/cluster/ClusterGroupAdapter.java
+++ b/modules/core/src/main/java/org/apache/ignite/internal/cluster/ClusterGroupAdapter.java
@@ -274,7 +274,7 @@ public ExecutorService executorService() {
if (nodes().isEmpty())
throw U.convertException(U.emptyTopologyException());
- return new ClusterMetricsSnapshot(this);
+ return new ClusterMetricsSnapshot(nodes());
}
finally {
unguard();
diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cluster/NodeFullMetricsMessage.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cluster/NodeFullMetricsMessage.java
index b52a355b61083..3902e05e2f454 100644
--- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cluster/NodeFullMetricsMessage.java
+++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cluster/NodeFullMetricsMessage.java
@@ -20,6 +20,7 @@
import java.util.Map;
import org.apache.ignite.cache.CacheMetrics;
import org.apache.ignite.cluster.ClusterMetrics;
+import org.apache.ignite.internal.ClusterMetricsSnapshot;
import org.apache.ignite.internal.Order;
import org.apache.ignite.internal.util.typedef.internal.S;
import org.apache.ignite.internal.util.typedef.internal.U;
@@ -30,7 +31,7 @@
public class NodeFullMetricsMessage implements Message {
/** Node metrics wrapper message. */
@Order(0)
- public NodeMetricsMessage nodeMetricsMsg;
+ public ClusterMetricsSnapshot nodeMetricsMsg;
/** Cache metrics wrapper message. */
@Order(1)
@@ -43,7 +44,7 @@ public NodeFullMetricsMessage() {
/** */
public NodeFullMetricsMessage(ClusterMetrics nodeMetrics, Map cacheMetrics) {
- nodeMetricsMsg = new NodeMetricsMessage(nodeMetrics);
+ nodeMetricsMsg = new ClusterMetricsSnapshot(nodeMetrics);
cachesMetricsMsgs = U.newHashMap(cacheMetrics.size());
@@ -61,12 +62,12 @@ public void cachesMetricsMessages(Map cacheMetrics
}
/** */
- public NodeMetricsMessage nodeMetricsMessage() {
+ public ClusterMetricsSnapshot nodeMetricsMessage() {
return nodeMetricsMsg;
}
/** */
- public void nodeMetricsMessage(NodeMetricsMessage nodeMetricsMsg) {
+ public void nodeMetricsMessage(ClusterMetricsSnapshot nodeMetricsMsg) {
this.nodeMetricsMsg = nodeMetricsMsg;
}
diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cluster/NodeMetricsMessage.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cluster/NodeMetricsMessage.java
deleted file mode 100644
index e97b5743434c5..0000000000000
--- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cluster/NodeMetricsMessage.java
+++ /dev/null
@@ -1,1344 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.ignite.internal.processors.cluster;
-
-import java.util.Collection;
-import java.util.Map;
-import org.apache.ignite.cluster.ClusterMetrics;
-import org.apache.ignite.cluster.ClusterNode;
-import org.apache.ignite.internal.Order;
-import org.apache.ignite.internal.util.typedef.F;
-import org.apache.ignite.internal.util.typedef.internal.S;
-import org.apache.ignite.internal.util.typedef.internal.U;
-import org.apache.ignite.plugin.extensions.communication.Message;
-import org.jetbrains.annotations.Nullable;
-
-import static java.lang.Math.max;
-import static java.lang.Math.min;
-
-/** */
-public class NodeMetricsMessage implements Message {
- /** */
- @Order(0)
- public long lastUpdateTime = -1;
-
- /** */
- @Order(1)
- public int maxActiveJobs = -1;
-
- /** */
- @Order(2)
- public int curActiveJobs = -1;
-
- /** */
- @Order(3)
- public float avgActiveJobs = -1;
-
- /** */
- @Order(4)
- public int maxWaitingJobs = -1;
-
- /** */
- @Order(5)
- public int curWaitingJobs = -1;
-
- /** */
- @Order(6)
- public float avgWaitingJobs = -1;
-
- /** */
- @Order(7)
- public int maxRejectedJobs = -1;
-
- /** */
- @Order(8)
- public int curRejectedJobs = -1;
-
- /** */
- @Order(9)
- public float avgRejectedJobs = -1;
-
- /** */
- @Order(10)
- public int maxCancelledJobs = -1;
-
- /** */
- @Order(11)
- public int curCancelledJobs = -1;
-
- /** */
- @Order(12)
- public float avgCancelledJobs = -1;
-
- /** */
- @Order(13)
- public int totalRejectedJobs = -1;
-
- /** */
- @Order(14)
- public int totalCancelledJobs = -1;
-
- /** */
- @Order(15)
- public int totalExecutedJobs = -1;
-
- /** */
- @Order(16)
- public long maxJobWaitTime = -1;
-
- /** */
- @Order(17)
- public long curJobWaitTime = Long.MAX_VALUE;
-
- /** */
- @Order(18)
- public double avgJobWaitTime = -1;
-
- /** */
- @Order(19)
- public long maxJobExecTime = -1;
-
- /** */
- @Order(20)
- public long curJobExecTime = -1;
-
- /** */
- @Order(21)
- public double avgJobExecTime = -1;
-
- /** */
- @Order(22)
- public int totalExecTasks = -1;
-
- /** */
- @Order(23)
- public long totalIdleTime = -1;
-
- /** */
- @Order(24)
- public long curIdleTime = -1;
-
- /** */
- @Order(25)
- public int totalCpus = -1;
-
- /** */
- @Order(26)
- public double curCpuLoad = -1;
-
- /** */
- @Order(27)
- public double avgCpuLoad = -1;
-
- /** */
- @Order(28)
- public double curGcCpuLoad = -1;
-
- /** */
- @Order(29)
- public long heapInit = -1;
-
- /** */
- @Order(30)
- public long heapUsed = -1;
-
- /** */
- @Order(31)
- public long heapCommitted = -1;
-
- /** */
- @Order(32)
- public long heapMax = -1;
-
- /** */
- @Order(33)
- public long heapTotal = -1;
-
- /** */
- @Order(34)
- public long nonHeapInit = -1;
-
- /** */
- @Order(35)
- public long nonHeapUsed = -1;
-
- /** */
- @Order(36)
- public long nonHeapCommitted = -1;
-
- /** */
- @Order(37)
- public long nonHeapMax = -1;
-
- /** */
- @Order(38)
- public long nonHeapTotal = -1;
-
- /** */
- @Order(39)
- public long upTime = -1;
-
- /** */
- @Order(40)
- public long startTime = -1;
-
- /** */
- @Order(41)
- public long nodeStartTime = -1;
-
- /** */
- @Order(42)
- public int threadCnt = -1;
-
- /** */
- @Order(43)
- public int peakThreadCnt = -1;
-
- /** */
- @Order(44)
- public long startedThreadCnt = -1;
-
- /** */
- @Order(45)
- public int daemonThreadCnt = -1;
-
- /** */
- @Order(46)
- public long lastDataVer = -1;
-
- /** */
- @Order(47)
- public int sentMsgsCnt = -1;
-
- /** */
- @Order(48)
- public long sentBytesCnt = -1;
-
- /** */
- @Order(49)
- public int rcvdMsgsCnt = -1;
-
- /** */
- @Order(50)
- public long rcvdBytesCnt = -1;
-
- /** */
- @Order(51)
- public int outMesQueueSize = -1;
-
- /** */
- @Order(52)
- public int totalNodes = -1;
-
- /** */
- @Order(53)
- public long totalJobsExecTime = -1;
-
- /** */
- @Order(54)
- public long curPmeDuration = -1;
-
- /** */
- public NodeMetricsMessage() {
- // No-op.
- }
-
- /**
- * Create metrics for given nodes.
- *
- * @param nodes Nodes.
- */
- public NodeMetricsMessage(Collection nodes) {
- int size = nodes.size();
-
- curJobWaitTime = Long.MAX_VALUE;
- lastUpdateTime = 0;
- maxActiveJobs = 0;
- curActiveJobs = 0;
- avgActiveJobs = 0;
- maxWaitingJobs = 0;
- curWaitingJobs = 0;
- avgWaitingJobs = 0;
- maxRejectedJobs = 0;
- curRejectedJobs = 0;
- avgRejectedJobs = 0;
- maxCancelledJobs = 0;
- curCancelledJobs = 0;
- avgCancelledJobs = 0;
- totalRejectedJobs = 0;
- totalCancelledJobs = 0;
- totalExecutedJobs = 0;
- totalJobsExecTime = 0;
- maxJobWaitTime = 0;
- avgJobWaitTime = 0;
- maxJobExecTime = 0;
- curJobExecTime = 0;
- avgJobExecTime = 0;
- totalExecTasks = 0;
- totalIdleTime = 0;
- curIdleTime = 0;
- totalCpus = 0;
- curCpuLoad = 0;
- avgCpuLoad = 0;
- curGcCpuLoad = 0;
- heapInit = 0;
- heapUsed = 0;
- heapCommitted = 0;
- heapMax = 0;
- nonHeapInit = 0;
- nonHeapUsed = 0;
- nonHeapCommitted = 0;
- nonHeapMax = 0;
- nonHeapTotal = 0;
- upTime = 0;
- startTime = 0;
- nodeStartTime = 0;
- threadCnt = 0;
- peakThreadCnt = 0;
- startedThreadCnt = 0;
- daemonThreadCnt = 0;
- lastDataVer = 0;
- sentMsgsCnt = 0;
- sentBytesCnt = 0;
- rcvdMsgsCnt = 0;
- rcvdBytesCnt = 0;
- outMesQueueSize = 0;
- heapTotal = 0;
- totalNodes = nodes.size();
- curPmeDuration = 0;
-
- for (ClusterNode node : nodes) {
- ClusterMetrics m = node.metrics();
-
- lastUpdateTime = max(lastUpdateTime, node.metrics().getLastUpdateTime());
-
- curActiveJobs += m.getCurrentActiveJobs();
- maxActiveJobs = max(maxActiveJobs, m.getCurrentActiveJobs());
- avgActiveJobs += m.getCurrentActiveJobs();
- totalExecutedJobs += m.getTotalExecutedJobs();
- totalJobsExecTime += m.getTotalJobsExecutionTime();
-
- totalExecTasks += m.getTotalExecutedTasks();
-
- totalCancelledJobs += m.getTotalCancelledJobs();
- curCancelledJobs += m.getCurrentCancelledJobs();
- maxCancelledJobs = max(maxCancelledJobs, m.getCurrentCancelledJobs());
- avgCancelledJobs += m.getCurrentCancelledJobs();
-
- totalRejectedJobs += m.getTotalRejectedJobs();
- curRejectedJobs += m.getCurrentRejectedJobs();
- maxRejectedJobs = max(maxRejectedJobs, m.getCurrentRejectedJobs());
- avgRejectedJobs += m.getCurrentRejectedJobs();
-
- curWaitingJobs += m.getCurrentWaitingJobs();
- maxWaitingJobs = max(maxWaitingJobs, m.getCurrentWaitingJobs());
- avgWaitingJobs += m.getCurrentWaitingJobs();
-
- maxJobExecTime = max(maxJobExecTime, m.getMaximumJobExecuteTime());
- avgJobExecTime += m.getAverageJobExecuteTime();
- curJobExecTime += m.getCurrentJobExecuteTime();
-
- curJobWaitTime = min(curJobWaitTime, m.getCurrentJobWaitTime());
- maxJobWaitTime = max(maxJobWaitTime, m.getCurrentJobWaitTime());
- avgJobWaitTime += m.getAverageJobWaitTime();
-
- daemonThreadCnt += m.getCurrentDaemonThreadCount();
-
- peakThreadCnt = max(peakThreadCnt, m.getCurrentThreadCount());
- threadCnt += m.getCurrentThreadCount();
- startedThreadCnt += m.getTotalStartedThreadCount();
-
- curIdleTime += m.getCurrentIdleTime();
- totalIdleTime += m.getTotalIdleTime();
-
- heapCommitted += m.getHeapMemoryCommitted();
-
- heapUsed += m.getHeapMemoryUsed();
-
- heapMax = max(heapMax, m.getHeapMemoryMaximum());
-
- heapTotal += m.getHeapMemoryTotal();
-
- heapInit += m.getHeapMemoryInitialized();
-
- nonHeapCommitted += m.getNonHeapMemoryCommitted();
-
- nonHeapUsed += m.getNonHeapMemoryUsed();
-
- nonHeapMax = max(nonHeapMax, m.getNonHeapMemoryMaximum());
-
- nonHeapTotal += m.getNonHeapMemoryTotal();
-
- nonHeapInit += m.getNonHeapMemoryInitialized();
-
- upTime = max(upTime, m.getUpTime());
-
- lastDataVer = max(lastDataVer, m.getLastDataVersion());
-
- sentMsgsCnt += m.getSentMessagesCount();
- sentBytesCnt += m.getSentBytesCount();
- rcvdMsgsCnt += m.getReceivedMessagesCount();
- rcvdBytesCnt += m.getReceivedBytesCount();
- outMesQueueSize += m.getOutboundMessagesQueueSize();
-
- avgCpuLoad += m.getCurrentCpuLoad();
-
- curPmeDuration = max(curPmeDuration, m.getCurrentPmeDuration());
- }
-
- curJobExecTime /= size;
-
- avgActiveJobs /= size;
- avgCancelledJobs /= size;
- avgRejectedJobs /= size;
- avgWaitingJobs /= size;
- avgJobExecTime /= size;
- avgJobWaitTime /= size;
- avgCpuLoad /= size;
-
- if (!F.isEmpty(nodes)) {
- ClusterMetrics oldestNodeMetrics = oldest(nodes).metrics();
-
- nodeStartTime = oldestNodeMetrics.getNodeStartTime();
- startTime = oldestNodeMetrics.getStartTime();
- }
-
- Map> neighborhood = U.neighborhood(nodes);
-
- curGcCpuLoad = currentGcCpuLoad(neighborhood);
- curCpuLoad = currentCpuLoad(neighborhood);
- totalCpus = cpuCnt(neighborhood);
- }
-
- /** */
- public NodeMetricsMessage(ClusterMetrics metrics) {
- maxActiveJobs = metrics.getMaximumActiveJobs();
- curActiveJobs = metrics.getCurrentActiveJobs();
- avgActiveJobs = metrics.getAverageActiveJobs();
-
- maxWaitingJobs = metrics.getMaximumWaitingJobs();
- curWaitingJobs = metrics.getCurrentWaitingJobs();
- avgWaitingJobs = metrics.getAverageWaitingJobs();
-
- maxRejectedJobs = metrics.getMaximumRejectedJobs();
- curRejectedJobs = metrics.getCurrentRejectedJobs();
- avgRejectedJobs = metrics.getAverageRejectedJobs();
-
- maxCancelledJobs = metrics.getMaximumCancelledJobs();
- curCancelledJobs = metrics.getCurrentCancelledJobs();
- avgCancelledJobs = metrics.getAverageCancelledJobs();
-
- totalRejectedJobs = metrics.getTotalRejectedJobs();
- totalCancelledJobs = metrics.getTotalCancelledJobs();
- totalExecutedJobs = metrics.getTotalExecutedJobs();
-
- maxJobWaitTime = metrics.getMaximumJobWaitTime();
- curJobWaitTime = metrics.getCurrentJobWaitTime();
- avgJobWaitTime = metrics.getAverageJobWaitTime();
-
- maxJobExecTime = metrics.getMaximumJobExecuteTime();
- curJobExecTime = metrics.getCurrentJobExecuteTime();
- avgJobExecTime = metrics.getAverageJobExecuteTime();
-
- totalJobsExecTime = metrics.getTotalJobsExecutionTime();
- totalExecTasks = metrics.getTotalExecutedTasks();
-
- curIdleTime = metrics.getCurrentIdleTime();
- totalIdleTime = metrics.getTotalIdleTime();
-
- totalCpus = metrics.getTotalCpus();
- curCpuLoad = metrics.getCurrentCpuLoad();
- avgCpuLoad = metrics.getAverageCpuLoad();
- curGcCpuLoad = metrics.getCurrentGcCpuLoad();
-
- heapInit = metrics.getHeapMemoryInitialized();
- heapUsed = metrics.getHeapMemoryUsed();
- heapCommitted = metrics.getHeapMemoryCommitted();
- heapMax = metrics.getHeapMemoryMaximum();
- heapTotal = metrics.getHeapMemoryTotal();
-
- nonHeapInit = metrics.getNonHeapMemoryInitialized();
- nonHeapUsed = metrics.getNonHeapMemoryUsed();
- nonHeapCommitted = metrics.getNonHeapMemoryCommitted();
- nonHeapMax = metrics.getNonHeapMemoryMaximum();
- nonHeapTotal = metrics.getNonHeapMemoryTotal();
-
- startTime = metrics.getStartTime();
- nodeStartTime = metrics.getNodeStartTime();
- upTime = metrics.getUpTime();
-
- lastDataVer = metrics.getLastDataVersion();
-
- curPmeDuration = metrics.getCurrentPmeDuration();
-
- totalNodes = metrics.getTotalNodes();
-
- threadCnt = metrics.getCurrentThreadCount();
- peakThreadCnt = metrics.getMaximumThreadCount();
- startedThreadCnt = metrics.getTotalStartedThreadCount();
- daemonThreadCnt = metrics.getCurrentDaemonThreadCount();
-
- sentMsgsCnt = metrics.getSentMessagesCount();
- rcvdMsgsCnt = metrics.getReceivedMessagesCount();
- outMesQueueSize = metrics.getOutboundMessagesQueueSize();
-
- sentBytesCnt = metrics.getSentBytesCount();
- rcvdBytesCnt = metrics.getReceivedBytesCount();
- }
-
- /** */
- public long heapMemoryTotal() {
- return heapTotal;
- }
-
- /**
- * Sets total heap size.
- *
- * @param heapTotal Total heap.
- */
- public void heapMemoryTotal(long heapTotal) {
- this.heapTotal = heapTotal;
- }
-
- /**
- * Sets non-heap total heap size.
- *
- * @param nonHeapTotal Total heap.
- */
- public void nonHeapMemoryTotal(long nonHeapTotal) {
- this.nonHeapTotal = nonHeapTotal;
- }
-
- /** */
- public long lastUpdateTime() {
- return lastUpdateTime;
- }
-
- /**
- * Sets last update time.
- *
- * @param lastUpdateTime Last update time.
- */
- public void lastUpdateTime(long lastUpdateTime) {
- this.lastUpdateTime = lastUpdateTime;
- }
-
- /** */
- public int maximumActiveJobs() {
- return maxActiveJobs;
- }
-
- /**
- * Sets max active jobs.
- *
- * @param maxActiveJobs Max active jobs.
- */
- public void maximumActiveJobs(int maxActiveJobs) {
- this.maxActiveJobs = maxActiveJobs;
- }
-
- /** */
- public int currentActiveJobs() {
- return curActiveJobs;
- }
-
- /**
- * Sets current active jobs.
- *
- * @param curActiveJobs Current active jobs.
- */
- public void currentActiveJobs(int curActiveJobs) {
- this.curActiveJobs = curActiveJobs;
- }
-
- /** */
- public float averageActiveJobs() {
- return avgActiveJobs;
- }
-
- /**
- * Sets average active jobs.
- *
- * @param avgActiveJobs Average active jobs.
- */
- public void averageActiveJobs(float avgActiveJobs) {
- this.avgActiveJobs = avgActiveJobs;
- }
-
- /** */
- public int maximumWaitingJobs() {
- return maxWaitingJobs;
- }
-
- /**
- * Sets maximum waiting jobs.
- *
- * @param maxWaitingJobs Maximum waiting jobs.
- */
- public void maximumWaitingJobs(int maxWaitingJobs) {
- this.maxWaitingJobs = maxWaitingJobs;
- }
-
- /** */
- public int currentWaitingJobs() {
- return curWaitingJobs;
- }
-
- /**
- * Sets current waiting jobs.
- *
- * @param curWaitingJobs Current waiting jobs.
- */
- public void currentWaitingJobs(int curWaitingJobs) {
- this.curWaitingJobs = curWaitingJobs;
- }
-
- /** */
- public float averageWaitingJobs() {
- return avgWaitingJobs;
- }
-
- /**
- * Sets average waiting jobs.
- *
- * @param avgWaitingJobs Average waiting jobs.
- */
- public void averageWaitingJobs(float avgWaitingJobs) {
- this.avgWaitingJobs = avgWaitingJobs;
- }
-
- /** */
- public int maximumRejectedJobs() {
- return maxRejectedJobs;
- }
-
- /**
- * @param maxRejectedJobs Maximum number of jobs rejected during a single collision resolution event.
- */
- public void maximumRejectedJobs(int maxRejectedJobs) {
- this.maxRejectedJobs = maxRejectedJobs;
- }
-
- /** */
- public int currentRejectedJobs() {
- return curRejectedJobs;
- }
-
- /**
- * @param curRejectedJobs Number of jobs rejected during most recent collision resolution.
- */
- public void currentRejectedJobs(int curRejectedJobs) {
- this.curRejectedJobs = curRejectedJobs;
- }
-
- /** */
- public float averageRejectedJobs() {
- return avgRejectedJobs;
- }
-
- /**
- * @param avgRejectedJobs Average number of jobs this node rejects.
- */
- public void averageRejectedJobs(float avgRejectedJobs) {
- this.avgRejectedJobs = avgRejectedJobs;
- }
-
- /** */
- public int totalRejectedJobs() {
- return totalRejectedJobs;
- }
-
- /**
- * @param totalRejectedJobs Total number of jobs this node ever rejected.
- */
- public void totalRejectedJobs(int totalRejectedJobs) {
- this.totalRejectedJobs = totalRejectedJobs;
- }
-
- /** */
- public int maximumCancelledJobs() {
- return maxCancelledJobs;
- }
-
- /**
- * Sets maximum cancelled jobs.
- *
- * @param maxCancelledJobs Maximum cancelled jobs.
- */
- public void maximumCancelledJobs(int maxCancelledJobs) {
- this.maxCancelledJobs = maxCancelledJobs;
- }
-
- /** */
- public int currentCancelledJobs() {
- return curCancelledJobs;
- }
-
- /**
- * Sets current cancelled jobs.
- *
- * @param curCancelledJobs Current cancelled jobs.
- */
- public void currentCancelledJobs(int curCancelledJobs) {
- this.curCancelledJobs = curCancelledJobs;
- }
-
- /** */
- public float averageCancelledJobs() {
- return avgCancelledJobs;
- }
-
- /**
- * Sets average cancelled jobs.
- *
- * @param avgCancelledJobs Average cancelled jobs.
- */
- public void averageCancelledJobs(float avgCancelledJobs) {
- this.avgCancelledJobs = avgCancelledJobs;
- }
-
- /** */
- public int totalExecutedJobs() {
- return totalExecutedJobs;
- }
-
- /**
- * Sets total active jobs.
- *
- * @param totalExecutedJobs Total active jobs.
- */
- public void totalExecutedJobs(int totalExecutedJobs) {
- this.totalExecutedJobs = totalExecutedJobs;
- }
-
- /** */
- public long totalJobsExecutionTime() {
- return totalJobsExecTime;
- }
-
- /**
- * Sets total jobs execution time.
- *
- * @param totalJobsExecTime Total jobs execution time.
- */
- public void totalJobsExecutionTime(long totalJobsExecTime) {
- this.totalJobsExecTime = totalJobsExecTime;
- }
-
- /** */
- public int totalCancelledJobs() {
- return totalCancelledJobs;
- }
-
- /**
- * Sets total cancelled jobs.
- *
- * @param totalCancelledJobs Total cancelled jobs.
- */
- public void totalCancelledJobs(int totalCancelledJobs) {
- this.totalCancelledJobs = totalCancelledJobs;
- }
-
- /** */
- public long maximumJobWaitTime() {
- return maxJobWaitTime;
- }
-
- /**
- * Sets max job wait time.
- *
- * @param maxJobWaitTime Max job wait time.
- */
- public void maximumJobWaitTime(long maxJobWaitTime) {
- this.maxJobWaitTime = maxJobWaitTime;
- }
-
- /** */
- public long currentJobWaitTime() {
- return curJobWaitTime;
- }
-
- /**
- * Sets current job wait time.
- *
- * @param curJobWaitTime Current job wait time.
- */
- public void currentJobWaitTime(long curJobWaitTime) {
- this.curJobWaitTime = curJobWaitTime;
- }
-
- /** */
- public double averageJobWaitTime() {
- return avgJobWaitTime;
- }
-
- /**
- * Sets average job wait time.
- *
- * @param avgJobWaitTime Average job wait time.
- */
- public void averageJobWaitTime(double avgJobWaitTime) {
- this.avgJobWaitTime = avgJobWaitTime;
- }
-
- /** */
- public long maximumJobExecuteTime() {
- return maxJobExecTime;
- }
-
- /**
- * Sets maximum job execution time.
- *
- * @param maxJobExecTime Maximum job execution time.
- */
- public void maximumJobExecuteTime(long maxJobExecTime) {
- this.maxJobExecTime = maxJobExecTime;
- }
-
- /** */
- public long currentJobExecuteTime() {
- return curJobExecTime;
- }
-
- /**
- * Sets current job execute time.
- *
- * @param curJobExecTime Current job execute time.
- */
- public void currentJobExecuteTime(long curJobExecTime) {
- this.curJobExecTime = curJobExecTime;
- }
-
- /** */
- public double averageJobExecuteTime() {
- return avgJobExecTime;
- }
-
- /**
- * Sets average job execution time.
- *
- * @param avgJobExecTime Average job execution time.
- */
- public void averageJobExecuteTime(double avgJobExecTime) {
- this.avgJobExecTime = avgJobExecTime;
- }
-
- /** */
- public int totalExecutedTasks() {
- return totalExecTasks;
- }
-
- /**
- * Sets total executed tasks count.
- *
- * @param totalExecTasks total executed tasks count.
- */
- public void totalExecutedTasks(int totalExecTasks) {
- this.totalExecTasks = totalExecTasks;
- }
-
- /** */
- public long totalIdleTime() {
- return totalIdleTime;
- }
-
- /**
- * Set total node idle time.
- *
- * @param totalIdleTime Total node idle time.
- */
- public void totalIdleTime(long totalIdleTime) {
- this.totalIdleTime = totalIdleTime;
- }
-
- /** */
- public long currentIdleTime() {
- return curIdleTime;
- }
-
- /**
- * Sets time elapsed since execution of last job.
- *
- * @param curIdleTime Time elapsed since execution of last job.
- */
- public void currentIdleTime(long curIdleTime) {
- this.curIdleTime = curIdleTime;
- }
-
- /** */
- public int totalCpus() {
- return totalCpus;
- }
-
- /** */
- public double currentCpuLoad() {
- return curCpuLoad;
- }
-
- /** */
- public double averageCpuLoad() {
- return avgCpuLoad;
- }
-
- /** */
- public double currentGcCpuLoad() {
- return curGcCpuLoad;
- }
-
- /** */
- public long heapMemoryInitialized() {
- return heapInit;
- }
-
- /** */
- public long heapMemoryUsed() {
- return heapUsed;
- }
-
- /** */
- public long heapMemoryCommitted() {
- return heapCommitted;
- }
-
- /** */
- public long heapMemoryMaximum() {
- return heapMax;
- }
-
- /** */
- public long nonHeapMemoryInitialized() {
- return nonHeapInit;
- }
-
- /** */
- public long nonHeapMemoryUsed() {
- return nonHeapUsed;
- }
-
- /** */
- public long nonHeapMemoryCommitted() {
- return nonHeapCommitted;
- }
-
- /** */
- public long nonHeapMemoryMaximum() {
- return nonHeapMax;
- }
-
- /** */
- public long nonHeapMemoryTotal() {
- return nonHeapTotal;
- }
-
- /** */
- public long upTime() {
- return upTime;
- }
-
- /** */
- public long startTime() {
- return startTime;
- }
-
- /** */
- public long nodeStartTime() {
- return nodeStartTime;
- }
-
- /** */
- public int currentThreadCount() {
- return threadCnt;
- }
-
- /** */
- public int maximumThreadCount() {
- return peakThreadCnt;
- }
-
- /** */
- public long totalStartedThreadCount() {
- return startedThreadCnt;
- }
-
- /** */
- public int currentDaemonThreadCount() {
- return daemonThreadCnt;
- }
-
- /** */
- public long lastDataVersion() {
- return lastDataVer;
- }
-
- /** */
- public int sentMessagesCount() {
- return sentMsgsCnt;
- }
-
- /** */
- public long sentBytesCount() {
- return sentBytesCnt;
- }
-
- /** */
- public int receivedMessagesCount() {
- return rcvdMsgsCnt;
- }
-
- /** */
- public long receivedBytesCount() {
- return rcvdBytesCnt;
- }
-
- /** */
- public int outboundMessagesQueueSize() {
- return outMesQueueSize;
- }
-
- /** */
- public int totalNodes() {
- return totalNodes;
- }
-
- /** */
- public long currentPmeDuration() {
- return curPmeDuration;
- }
-
- /**
- * Sets available processors.
- *
- * @param totalCpus Available processors.
- */
- public void totalCpus(int totalCpus) {
- this.totalCpus = totalCpus;
- }
-
- /**
- * Sets current CPU load.
- *
- * @param curCpuLoad Current CPU load.
- */
- public void currentCpuLoad(double curCpuLoad) {
- this.curCpuLoad = curCpuLoad;
- }
-
- /**
- * Sets CPU load average over the metrics history.
- *
- * @param avgCpuLoad CPU load average.
- */
- public void averageCpuLoad(double avgCpuLoad) {
- this.avgCpuLoad = avgCpuLoad;
- }
-
- /**
- * Sets current GC load.
- *
- * @param curGcCpuLoad Current GC load.
- */
- public void currentGcCpuLoad(double curGcCpuLoad) {
- this.curGcCpuLoad = curGcCpuLoad;
- }
-
- /**
- * Sets heap initial memory.
- *
- * @param heapInit Heap initial memory.
- */
- public void heapMemoryInitialized(long heapInit) {
- this.heapInit = heapInit;
- }
-
- /**
- * Sets used heap memory.
- *
- * @param heapUsed Used heap memory.
- */
- public void heapMemoryUsed(long heapUsed) {
- this.heapUsed = heapUsed;
- }
-
- /**
- * Sets committed heap memory.
- *
- * @param heapCommitted Committed heap memory.
- */
- public void heapMemoryCommitted(long heapCommitted) {
- this.heapCommitted = heapCommitted;
- }
-
- /**
- * Sets maximum possible heap memory.
- *
- * @param heapMax Maximum possible heap memory.
- */
- public void heapMemoryMaximum(long heapMax) {
- this.heapMax = heapMax;
- }
-
- /**
- * Sets initial non-heap memory.
- *
- * @param nonHeapInit Initial non-heap memory.
- */
- public void nonHeapMemoryInitialized(long nonHeapInit) {
- this.nonHeapInit = nonHeapInit;
- }
-
- /**
- * Sets used non-heap memory.
- *
- * @param nonHeapUsed Used non-heap memory.
- */
- public void nonHeapMemoryUsed(long nonHeapUsed) {
- this.nonHeapUsed = nonHeapUsed;
- }
-
- /**
- * Sets committed non-heap memory.
- *
- * @param nonHeapCommitted Committed non-heap memory.
- */
- public void nonHeapMemoryCommitted(long nonHeapCommitted) {
- this.nonHeapCommitted = nonHeapCommitted;
- }
-
- /**
- * Sets maximum possible non-heap memory.
- *
- * @param nonHeapMax Maximum possible non-heap memory.
- */
- public void nonHeapMemoryMaximum(long nonHeapMax) {
- this.nonHeapMax = nonHeapMax;
- }
-
- /**
- * Sets VM up time.
- *
- * @param upTime VM up time.
- */
- public void upTime(long upTime) {
- this.upTime = upTime;
- }
-
- /**
- * Sets VM start time.
- *
- * @param startTime VM start time.
- */
- public void startTime(long startTime) {
- this.startTime = startTime;
- }
-
- /**
- * Sets node start time.
- *
- * @param nodeStartTime node start time.
- */
- public void nodeStartTime(long nodeStartTime) {
- this.nodeStartTime = nodeStartTime;
- }
-
- /**
- * Sets thread count.
- *
- * @param threadCnt Thread count.
- */
- public void currentThreadCount(int threadCnt) {
- this.threadCnt = threadCnt;
- }
-
- /**
- * Sets peak thread count.
- *
- * @param peakThreadCnt Peak thread count.
- */
- public void maximumThreadCount(int peakThreadCnt) {
- this.peakThreadCnt = peakThreadCnt;
- }
-
- /**
- * Sets started thread count.
- *
- * @param startedThreadCnt Started thread count.
- */
- public void totalStartedThreadCount(long startedThreadCnt) {
- this.startedThreadCnt = startedThreadCnt;
- }
-
- /**
- * Sets daemon thread count.
- *
- * @param daemonThreadCnt Daemon thread count.
- */
- public void currentDaemonThreadCount(int daemonThreadCnt) {
- this.daemonThreadCnt = daemonThreadCnt;
- }
-
- /**
- * Sets last data version.
- *
- * @param lastDataVer Last data version.
- */
- public void lastDataVersion(long lastDataVer) {
- this.lastDataVer = lastDataVer;
- }
-
- /**
- * Sets sent messages count.
- *
- * @param sentMsgsCnt Sent messages count.
- */
- public void sentMessagesCount(int sentMsgsCnt) {
- this.sentMsgsCnt = sentMsgsCnt;
- }
-
- /**
- * Sets sent bytes count.
- *
- * @param sentBytesCnt Sent bytes count.
- */
- public void sentBytesCount(long sentBytesCnt) {
- this.sentBytesCnt = sentBytesCnt;
- }
-
- /**
- * Sets received messages count.
- *
- * @param rcvdMsgsCnt Received messages count.
- */
- public void receivedMessagesCount(int rcvdMsgsCnt) {
- this.rcvdMsgsCnt = rcvdMsgsCnt;
- }
-
- /**
- * Sets received bytes count.
- *
- * @param rcvdBytesCnt Received bytes count.
- */
- public void receivedBytesCount(long rcvdBytesCnt) {
- this.rcvdBytesCnt = rcvdBytesCnt;
- }
-
- /**
- * Sets outbound messages queue size.
- *
- * @param outMesQueueSize Outbound messages queue size.
- */
- public void outboundMessagesQueueSize(int outMesQueueSize) {
- this.outMesQueueSize = outMesQueueSize;
- }
-
- /**
- * Sets total number of nodes.
- *
- * @param totalNodes Total number of nodes.
- */
- public void totalNodes(int totalNodes) {
- this.totalNodes = totalNodes;
- }
-
- /**
- * Sets execution duration for current partition map exchange.
- *
- * @param curPmeDuration Execution duration for current partition map exchange.
- */
- public void currentPmeDuration(long curPmeDuration) {
- this.curPmeDuration = curPmeDuration;
- }
-
- /**
- * Gets the oldest node in given collection.
- *
- * @param nodes Nodes.
- * @return Oldest node or {@code null} if collection is empty.
- */
- @Nullable private static ClusterNode oldest(Collection nodes) {
- long min = Long.MAX_VALUE;
-
- ClusterNode oldest = null;
-
- for (ClusterNode n : nodes)
- if (n.order() < min) {
- min = n.order();
- oldest = n;
- }
-
- return oldest;
- }
-
- /**
- * @param neighborhood Cluster neighborhood.
- * @return CPU count.
- */
- private static int cpuCnt(Map> neighborhood) {
- int cpus = 0;
-
- for (Collection nodes : neighborhood.values()) {
- ClusterNode first = F.first(nodes);
-
- // Projection can be empty if all nodes in it failed.
- if (first != null)
- cpus += first.metrics().getTotalCpus();
- }
-
- return cpus;
- }
-
- /**
- * @param neighborhood Cluster neighborhood.
- * @return CPU load.
- */
- private static double currentCpuLoad(Map> neighborhood) {
- double curCpuLoad = 0.0;
-
- for (Collection nodes : neighborhood.values()) {
- ClusterNode first = F.first(nodes);
-
- // Projection can be empty if all nodes in it failed.
- if (first != null)
- curCpuLoad += first.metrics().getCurrentCpuLoad();
- }
-
- return curCpuLoad;
- }
-
- /**
- * @param neighborhood Cluster neighborhood.
- * @return GC CPU load.
- */
- private static double currentGcCpuLoad(Map> neighborhood) {
- double curGcCpuLoad = 0;
-
- for (Collection nodes : neighborhood.values()) {
- ClusterNode first = F.first(nodes);
-
- // Projection can be empty if all nodes in it failed.
- if (first != null)
- curGcCpuLoad += first.metrics().getCurrentGcCpuLoad();
- }
-
- return curGcCpuLoad;
- }
-
- /** */
- public String toString() {
- return S.toString(NodeMetricsMessage.class, this);
- }
-}
diff --git a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java
index 43d4f1b4867a3..5cdcdeb87f750 100644
--- a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java
+++ b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java
@@ -761,6 +761,10 @@ private static void sleepEx(long millis, Runnable before, Runnable after) throws
discoveryData = spi.collectExchangeData(dataPacket);
}
+ // Set initial metrics for node validation.
+ // TODO : Revise in https://issues.apache.org/jira/browse/IGNITE-28965
+ node.setMetrics(spi.metricsProvider.metrics());
+
msg = new TcpDiscoveryJoinRequestMessage(node, discoveryData);
}
else
diff --git a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java
index 1d70baed0aa48..5deb4e7fa3a87 100644
--- a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java
+++ b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java
@@ -1158,6 +1158,10 @@ private void joinTopology() throws IgniteSpiException {
DiscoveryDataPacket discoveryData = spi.collectExchangeData(new DiscoveryDataPacket(getLocalNodeId()));
+ // Set initial metrics for node validation.
+ // TODO : Revise in https://issues.apache.org/jira/browse/IGNITE-28965
+ locNode.setMetrics(spi.metricsProvider.metrics());
+
TcpDiscoveryJoinRequestMessage joinReqMsg = new TcpDiscoveryJoinRequestMessage(locNode, discoveryData);
while (true) {
diff --git a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryImpl.java b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryImpl.java
index 85af9fa6d4479..b83e2130ae2e0 100644
--- a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryImpl.java
+++ b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryImpl.java
@@ -40,7 +40,6 @@
import org.apache.ignite.internal.processors.cache.CacheMetricsSnapshot;
import org.apache.ignite.internal.processors.cluster.CacheMetricsMessage;
import org.apache.ignite.internal.processors.cluster.NodeFullMetricsMessage;
-import org.apache.ignite.internal.processors.cluster.NodeMetricsMessage;
import org.apache.ignite.internal.thread.context.OperationContextDispatcher;
import org.apache.ignite.internal.util.typedef.F;
import org.apache.ignite.internal.util.typedef.internal.LT;
@@ -408,7 +407,7 @@ public void processCacheMetricsMessage(TcpDiscoveryMetricsUpdateMessage msg, lon
for (Map.Entry e : msg.serversFullMetricsMessages().entrySet()) {
UUID srvrId = e.getKey();
Map cacheMetricsMsgs = e.getValue().cachesMetricsMessages();
- NodeMetricsMessage srvrMetricsMsg = e.getValue().nodeMetricsMessage();
+ ClusterMetricsSnapshot srvrMetricsMsg = e.getValue().nodeMetricsMessage();
assert srvrMetricsMsg != null;
diff --git a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/internal/TcpDiscoveryNode.java b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/internal/TcpDiscoveryNode.java
index 7c7f795b91883..8f1bee5ceb209 100644
--- a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/internal/TcpDiscoveryNode.java
+++ b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/internal/TcpDiscoveryNode.java
@@ -37,9 +37,7 @@
import org.apache.ignite.internal.IgniteNodeAttributes;
import org.apache.ignite.internal.Marshalled;
import org.apache.ignite.internal.Order;
-import org.apache.ignite.internal.SelfMarshallingMessage;
import org.apache.ignite.internal.managers.discovery.IgniteClusterNode;
-import org.apache.ignite.internal.processors.cluster.NodeMetricsMessage;
import org.apache.ignite.internal.util.lang.GridMetadataAwareAdapter;
import org.apache.ignite.internal.util.tostring.GridToStringExclude;
import org.apache.ignite.internal.util.tostring.GridToStringInclude;
@@ -48,6 +46,7 @@
import org.apache.ignite.internal.util.typedef.internal.U;
import org.apache.ignite.lang.IgnitePredicate;
import org.apache.ignite.lang.IgniteProductVersion;
+import org.apache.ignite.plugin.extensions.communication.Message;
import org.apache.ignite.spi.discovery.DiscoveryMetricsProvider;
import org.apache.ignite.spi.discovery.tcp.TcpDiscoverySpi;
import org.jetbrains.annotations.Nullable;
@@ -62,7 +61,7 @@
* public due to certain limitations of Java technology.
*/
public class TcpDiscoveryNode extends GridMetadataAwareAdapter implements IgniteClusterNode,
- Comparable, Externalizable, SelfMarshallingMessage {
+ Comparable, Externalizable, Message {
/** */
private static final long serialVersionUID = 0L;
@@ -108,16 +107,12 @@ public class TcpDiscoveryNode extends GridMetadataAwareAdapter implements Ignite
/** Node metrics. */
@GridToStringExclude
- volatile ClusterMetrics metrics;
-
- /** Node metrics message. */
- @GridToStringExclude
@Order(6)
- volatile NodeMetricsMessage metricsMsg;
+ volatile ClusterMetricsSnapshot clusterMetricsSnapshot;
/** Node cache metrics. */
@GridToStringExclude
- private volatile Map cacheMetrics;
+ private volatile Map cacheMetricsSnapshot;
/** Node order in the topology. */
@Order(7)
@@ -215,24 +210,9 @@ public TcpDiscoveryNode(UUID id,
this.consistentId = consistentId != null ? consistentId : U.consistentId(sortedAddrs, discPort);
- metrics = metricsProvider.metrics();
- cacheMetrics = metricsProvider.cacheMetrics();
sockAddrs = U.toSocketAddresses(this, discPort);
}
- /** {@inheritDoc} */
- @Override public void selfMarshal() {
- metricsMsg = new NodeMetricsMessage(metrics);
- }
-
- /** {@inheritDoc} */
- @Override public void selfUnmarshal() {
- if (metricsMsg != null)
- metrics = new ClusterMetricsSnapshot(metricsMsg);
-
- metricsMsg = null;
- }
-
/**
* @return Last successfully connected address.
*/
@@ -315,40 +295,31 @@ public Map getAttributes() {
/** {@inheritDoc} */
@Override public ClusterMetrics metrics() {
- if (metricsProvider != null) {
- ClusterMetrics metrics0 = metricsProvider.metrics();
+ assert clusterMetricsSnapshot != null || metricsProvider != null;
- metrics = metrics0;
-
- return metrics0;
- }
-
- return metrics;
+ return metricsProvider == null ? clusterMetricsSnapshot : metricsProvider.metrics();
}
/** {@inheritDoc} */
@Override public void setMetrics(ClusterMetrics metrics) {
assert metrics != null;
- this.metrics = metrics;
+ this.clusterMetricsSnapshot = ClusterMetricsSnapshot.of(metrics);
}
/** {@inheritDoc} */
@Override public Map cacheMetrics() {
if (metricsProvider != null) {
- Map cacheMetrics0 = metricsProvider.cacheMetrics();
-
- cacheMetrics = cacheMetrics0;
-
- return cacheMetrics0;
+ // TODO : Revise in https://issues.apache.org/jira/browse/IGNITE-28965
+ cacheMetricsSnapshot = metricsProvider.cacheMetrics();
}
- return cacheMetrics;
+ return cacheMetricsSnapshot;
}
/** {@inheritDoc} */
@Override public void setCacheMetrics(Map cacheMetrics) {
- this.cacheMetrics = cacheMetrics != null ? cacheMetrics : Collections.emptyMap();
+ this.cacheMetricsSnapshot = cacheMetrics != null ? cacheMetrics : Collections.emptyMap();
}
/**
@@ -609,7 +580,7 @@ public TcpDiscoveryNode clientReconnectNode(Map nodeAttrs) {
// Cluster metrics
byte[] mtr = null;
- ClusterMetrics metrics = this.metrics;
+ var metrics = this.clusterMetricsSnapshot;
if (metrics != null)
mtr = ClusterMetricsSnapshot.serialize(metrics);
@@ -640,7 +611,7 @@ public TcpDiscoveryNode clientReconnectNode(Map nodeAttrs) {
byte[] mtr = U.readByteArray(in);
if (mtr != null)
- metrics = ClusterMetricsSnapshot.deserialize(mtr, 0);
+ clusterMetricsSnapshot = ClusterMetricsSnapshot.deserialize(mtr, 0);
// Legacy: Cache metrics
int size = in.readInt();
diff --git a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/messages/TcpDiscoveryClientMetricsUpdateMessage.java b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/messages/TcpDiscoveryClientMetricsUpdateMessage.java
index 22e98a087a12e..1fbc24fa20fc9 100644
--- a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/messages/TcpDiscoveryClientMetricsUpdateMessage.java
+++ b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/messages/TcpDiscoveryClientMetricsUpdateMessage.java
@@ -19,8 +19,8 @@
import java.util.UUID;
import org.apache.ignite.cluster.ClusterMetrics;
+import org.apache.ignite.internal.ClusterMetricsSnapshot;
import org.apache.ignite.internal.Order;
-import org.apache.ignite.internal.processors.cluster.NodeMetricsMessage;
import org.apache.ignite.internal.util.typedef.internal.S;
import org.apache.ignite.plugin.extensions.communication.MessageFactory;
@@ -32,7 +32,7 @@
public class TcpDiscoveryClientMetricsUpdateMessage extends TcpDiscoveryAbstractMessage {
/** */
@Order(0)
- NodeMetricsMessage metricsMsg;
+ ClusterMetricsSnapshot metricsMsg;
/** Constructor for {@link MessageFactory}. */
public TcpDiscoveryClientMetricsUpdateMessage() {
@@ -48,7 +48,7 @@ public TcpDiscoveryClientMetricsUpdateMessage() {
public TcpDiscoveryClientMetricsUpdateMessage(UUID creatorNodeId, ClusterMetrics metrics) {
super(creatorNodeId);
- metricsMsg = new NodeMetricsMessage(metrics);
+ metricsMsg = new ClusterMetricsSnapshot(metrics);
}
/**
@@ -56,7 +56,7 @@ public TcpDiscoveryClientMetricsUpdateMessage(UUID creatorNodeId, ClusterMetrics
*
* @return Metrics holder message.
*/
- public NodeMetricsMessage metricsMessage() {
+ public ClusterMetricsSnapshot metricsMessage() {
return metricsMsg;
}
diff --git a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/messages/TcpDiscoveryClientNodesMetricsMessage.java b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/messages/TcpDiscoveryClientNodesMetricsMessage.java
index 7aa4afa0b0be8..26123c91e4ac5 100644
--- a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/messages/TcpDiscoveryClientNodesMetricsMessage.java
+++ b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/messages/TcpDiscoveryClientNodesMetricsMessage.java
@@ -19,8 +19,8 @@
import java.util.Map;
import java.util.UUID;
+import org.apache.ignite.internal.ClusterMetricsSnapshot;
import org.apache.ignite.internal.Order;
-import org.apache.ignite.internal.processors.cluster.NodeMetricsMessage;
import org.apache.ignite.internal.util.typedef.internal.S;
import org.apache.ignite.plugin.extensions.communication.Message;
import org.apache.ignite.plugin.extensions.communication.MessageFactory;
@@ -29,7 +29,7 @@
public class TcpDiscoveryClientNodesMetricsMessage implements Message {
/** Map of nodes metrics messages per node id. */
@Order(0)
- Map nodesMetricsMsgs;
+ Map nodesMetricsMsgs;
/** Constructor for {@link MessageFactory}. */
public TcpDiscoveryClientNodesMetricsMessage() {
@@ -37,12 +37,12 @@ public TcpDiscoveryClientNodesMetricsMessage() {
}
/** @return Map of nodes metrics messages per node id. */
- public Map nodesMetricsMessages() {
+ public Map nodesMetricsMessages() {
return nodesMetricsMsgs;
}
/** @param nodesMetricsMsgs Map of nodes metrics messages per node id. */
- public void nodesMetricsMessages(Map nodesMetricsMsgs) {
+ public void nodesMetricsMessages(Map nodesMetricsMsgs) {
this.nodesMetricsMsgs = nodesMetricsMsgs;
}
diff --git a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/messages/TcpDiscoveryMetricsUpdateMessage.java b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/messages/TcpDiscoveryMetricsUpdateMessage.java
index 28735010366ea..1258931e343e0 100644
--- a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/messages/TcpDiscoveryMetricsUpdateMessage.java
+++ b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/messages/TcpDiscoveryMetricsUpdateMessage.java
@@ -24,10 +24,10 @@
import java.util.UUID;
import org.apache.ignite.cache.CacheMetrics;
import org.apache.ignite.cluster.ClusterMetrics;
+import org.apache.ignite.internal.ClusterMetricsSnapshot;
import org.apache.ignite.internal.Order;
import org.apache.ignite.internal.processors.cluster.CacheMetricsMessage;
import org.apache.ignite.internal.processors.cluster.NodeFullMetricsMessage;
-import org.apache.ignite.internal.processors.cluster.NodeMetricsMessage;
import org.apache.ignite.internal.util.tostring.GridToStringExclude;
import org.apache.ignite.internal.util.typedef.internal.S;
import org.apache.ignite.internal.util.typedef.internal.U;
@@ -97,7 +97,7 @@ public void addServerMetrics(UUID srvrId, ClusterMetrics newMetrics) {
if (srvrFullMetrics == null)
srvrFullMetrics = new NodeFullMetricsMessage();
- srvrFullMetrics.nodeMetricsMessage(new NodeMetricsMessage(newMetrics));
+ srvrFullMetrics.nodeMetricsMessage(new ClusterMetricsSnapshot(newMetrics));
return srvrFullMetrics;
});
@@ -157,7 +157,7 @@ public void addClientMetrics(UUID srvrId, UUID clientNodeId, ClusterMetrics clie
clientsMetricsMsg.nodesMetricsMessages(new HashMap<>());
}
- clientsMetricsMsg.nodesMetricsMessages().put(clientNodeId, new NodeMetricsMessage(clientMetrics));
+ clientsMetricsMsg.nodesMetricsMessages().put(clientNodeId, new ClusterMetricsSnapshot(clientMetrics));
return clientsMetricsMsg;
});
diff --git a/modules/core/src/test/java/org/apache/ignite/spi/collision/jobstealing/GridJobStealingCollisionSpiAttributesSelfTest.java b/modules/core/src/test/java/org/apache/ignite/spi/collision/jobstealing/GridJobStealingCollisionSpiAttributesSelfTest.java
index f009f34046848..3bc191a3dc18b 100644
--- a/modules/core/src/test/java/org/apache/ignite/spi/collision/jobstealing/GridJobStealingCollisionSpiAttributesSelfTest.java
+++ b/modules/core/src/test/java/org/apache/ignite/spi/collision/jobstealing/GridJobStealingCollisionSpiAttributesSelfTest.java
@@ -24,7 +24,6 @@
import java.util.UUID;
import org.apache.ignite.cluster.ClusterNode;
import org.apache.ignite.internal.ClusterMetricsSnapshot;
-import org.apache.ignite.internal.processors.cluster.NodeMetricsMessage;
import org.apache.ignite.internal.util.typedef.F;
import org.apache.ignite.internal.util.typedef.internal.U;
import org.apache.ignite.spi.collision.CollisionJobContext;
@@ -110,7 +109,7 @@ private void addSpiDependency(GridTestNode node) throws Exception {
rmtNode.setAttribute(U.spiAttribute(getSpi(), WAIT_JOBS_THRESHOLD_NODE_ATTR), getWaitJobsThreshold());
- NodeMetricsMessage metrics = new NodeMetricsMessage();
+ ClusterMetricsSnapshot metrics = new ClusterMetricsSnapshot();
metrics.currentWaitingJobs(2);
diff --git a/modules/core/src/test/java/org/apache/ignite/spi/collision/jobstealing/GridJobStealingCollisionSpiCustomTopologySelfTest.java b/modules/core/src/test/java/org/apache/ignite/spi/collision/jobstealing/GridJobStealingCollisionSpiCustomTopologySelfTest.java
index dc02da51fcb3a..6957031e0c9b4 100644
--- a/modules/core/src/test/java/org/apache/ignite/spi/collision/jobstealing/GridJobStealingCollisionSpiCustomTopologySelfTest.java
+++ b/modules/core/src/test/java/org/apache/ignite/spi/collision/jobstealing/GridJobStealingCollisionSpiCustomTopologySelfTest.java
@@ -25,7 +25,6 @@
import org.apache.ignite.GridTestTaskSession;
import org.apache.ignite.cluster.ClusterNode;
import org.apache.ignite.internal.ClusterMetricsSnapshot;
-import org.apache.ignite.internal.processors.cluster.NodeMetricsMessage;
import org.apache.ignite.internal.util.typedef.internal.U;
import org.apache.ignite.lang.IgniteUuid;
import org.apache.ignite.spi.collision.CollisionJobContext;
@@ -91,7 +90,7 @@ public int getActiveJobsThreshold() {
addSpiDependency(rmtNode1);
addSpiDependency(rmtNode2);
- NodeMetricsMessage metricsMsg = new NodeMetricsMessage();
+ ClusterMetricsSnapshot metricsMsg = new ClusterMetricsSnapshot();
metricsMsg.currentWaitingJobs(2);
diff --git a/modules/core/src/test/java/org/apache/ignite/spi/collision/jobstealing/GridJobStealingCollisionSpiSelfTest.java b/modules/core/src/test/java/org/apache/ignite/spi/collision/jobstealing/GridJobStealingCollisionSpiSelfTest.java
index 853a7a2c6bf8c..f143e7d038e3f 100644
--- a/modules/core/src/test/java/org/apache/ignite/spi/collision/jobstealing/GridJobStealingCollisionSpiSelfTest.java
+++ b/modules/core/src/test/java/org/apache/ignite/spi/collision/jobstealing/GridJobStealingCollisionSpiSelfTest.java
@@ -25,7 +25,6 @@
import org.apache.ignite.GridTestTaskSession;
import org.apache.ignite.cluster.ClusterNode;
import org.apache.ignite.internal.ClusterMetricsSnapshot;
-import org.apache.ignite.internal.processors.cluster.NodeMetricsMessage;
import org.apache.ignite.internal.util.typedef.CI1;
import org.apache.ignite.internal.util.typedef.F;
import org.apache.ignite.internal.util.typedef.internal.U;
@@ -98,7 +97,7 @@ public int getMaximumStealingAttempts() {
rmtNode.setAttribute(U.spiAttribute(getSpi(), WAIT_JOBS_THRESHOLD_NODE_ATTR), getWaitJobsThreshold());
- NodeMetricsMessage metrics = new NodeMetricsMessage();
+ ClusterMetricsSnapshot metrics = new ClusterMetricsSnapshot();
metrics.currentWaitingJobs(2);
diff --git a/modules/core/src/test/java/org/apache/ignite/spi/discovery/ClusterMetricsSnapshotSerializeSelfTest.java b/modules/core/src/test/java/org/apache/ignite/spi/discovery/ClusterMetricsSnapshotSerializeSelfTest.java
index f45010396c768..75945de92aea6 100644
--- a/modules/core/src/test/java/org/apache/ignite/spi/discovery/ClusterMetricsSnapshotSerializeSelfTest.java
+++ b/modules/core/src/test/java/org/apache/ignite/spi/discovery/ClusterMetricsSnapshotSerializeSelfTest.java
@@ -19,7 +19,6 @@
import org.apache.ignite.cluster.ClusterMetrics;
import org.apache.ignite.internal.ClusterMetricsSnapshot;
-import org.apache.ignite.internal.processors.cluster.NodeMetricsMessage;
import org.apache.ignite.testframework.junits.common.GridCommonAbstractTest;
import org.apache.ignite.testframework.junits.common.GridCommonTest;
import org.junit.Test;
@@ -98,7 +97,7 @@ public void testMetricsCompatibility() {
* @return Test metrics.
*/
private ClusterMetrics createMetrics() {
- NodeMetricsMessage metrics = new NodeMetricsMessage();
+ ClusterMetricsSnapshot metrics = new ClusterMetricsSnapshot();
metrics.totalCpus(1);
metrics.averageActiveJobs(2);
diff --git a/modules/core/src/test/java/org/apache/ignite/testframework/GridSpiTestContext.java b/modules/core/src/test/java/org/apache/ignite/testframework/GridSpiTestContext.java
index db82f7095ab21..4b513005fc49a 100644
--- a/modules/core/src/test/java/org/apache/ignite/testframework/GridSpiTestContext.java
+++ b/modules/core/src/test/java/org/apache/ignite/testframework/GridSpiTestContext.java
@@ -46,7 +46,6 @@
import org.apache.ignite.internal.managers.communication.GridMessageListener;
import org.apache.ignite.internal.managers.communication.IgniteMessageFactoryImpl;
import org.apache.ignite.internal.managers.eventstorage.GridLocalEventListener;
-import org.apache.ignite.internal.processors.cluster.NodeMetricsMessage;
import org.apache.ignite.internal.processors.metric.MetricRegistryImpl;
import org.apache.ignite.internal.processors.timeout.GridSpiTimeoutObject;
import org.apache.ignite.internal.processors.timeout.GridTimeoutProcessor;
@@ -192,7 +191,7 @@ public void reset() {
* @return Metrics adapter.
*/
private ClusterMetricsSnapshot createMetrics(int waitingJobs, int activeJobs) {
- NodeMetricsMessage metrics = new NodeMetricsMessage();
+ ClusterMetricsSnapshot metrics = new ClusterMetricsSnapshot();
metrics.currentWaitingJobs(waitingJobs);
metrics.currentActiveJobs(activeJobs);
diff --git a/modules/core/src/test/java/org/apache/ignite/util/GridTopologyHeapSizeSelfTest.java b/modules/core/src/test/java/org/apache/ignite/util/GridTopologyHeapSizeSelfTest.java
index 2123647b20476..8523d02446cba 100644
--- a/modules/core/src/test/java/org/apache/ignite/util/GridTopologyHeapSizeSelfTest.java
+++ b/modules/core/src/test/java/org/apache/ignite/util/GridTopologyHeapSizeSelfTest.java
@@ -20,7 +20,6 @@
import java.util.UUID;
import org.apache.ignite.cluster.ClusterNode;
import org.apache.ignite.internal.ClusterMetricsSnapshot;
-import org.apache.ignite.internal.processors.cluster.NodeMetricsMessage;
import org.apache.ignite.internal.util.typedef.F;
import org.apache.ignite.internal.util.typedef.internal.U;
import org.apache.ignite.testframework.GridTestNode;
@@ -90,7 +89,7 @@ public void testTopologyHeapSizeForNodesWithDifferentMacs() {
* @return Node.
*/
private GridTestNode getNode(String mac, int pid) {
- NodeMetricsMessage metrics = new NodeMetricsMessage();
+ ClusterMetricsSnapshot metrics = new ClusterMetricsSnapshot();
metrics.heapMemoryMaximum(1024L * 1024 * 1024);
metrics.heapMemoryInitialized(1024L * 1024 * 1024);
diff --git a/modules/indexing/src/test/java/org/apache/ignite/internal/processors/query/SqlSystemViewsSelfTest.java b/modules/indexing/src/test/java/org/apache/ignite/internal/processors/query/SqlSystemViewsSelfTest.java
index 87f000f273822..cbb7fe5e55efd 100644
--- a/modules/indexing/src/test/java/org/apache/ignite/internal/processors/query/SqlSystemViewsSelfTest.java
+++ b/modules/indexing/src/test/java/org/apache/ignite/internal/processors/query/SqlSystemViewsSelfTest.java
@@ -73,7 +73,6 @@
import org.apache.ignite.internal.IgniteInterruptedCheckedException;
import org.apache.ignite.internal.IgniteNodeAttributes;
import org.apache.ignite.internal.cache.query.index.IndexProcessor;
-import org.apache.ignite.internal.managers.discovery.ClusterMetricsImpl;
import org.apache.ignite.internal.processors.cache.GridCacheProcessor;
import org.apache.ignite.internal.processors.cache.index.AbstractIndexingCommonTest;
import org.apache.ignite.internal.processors.cache.index.AbstractSchemaSelfTest;
@@ -1809,9 +1808,9 @@ public void testDurationMetricsCanBeLonger24Hours() throws Exception {
// Get rid of metrics provider: current logic ignores metrics field if provider != null.
setField(node, "metricsProvider", null);
- ClusterMetricsImpl original = getField(node, "metrics");
+ ClusterMetricsSnapshot original = getField(node, "clusterMetricsSnapshot");
- setField(node, "metrics", new MockedClusterMetrics(original));
+ setField(node, "clusterMetricsSnapshot", new MockedClusterMetrics(original));
List> durationMetrics = execSql(ign,
"SELECT " +
@@ -1952,9 +1951,9 @@ public void testLocksView() throws Exception {
}
/**
- * Mock for {@link ClusterMetricsImpl} that always returns big (more than 24h) duration for all duration metrics.
+ * Mock for {@link ClusterMetricsSnapshot} that always returns big (more than 24h) duration for all duration metrics.
*/
- public static class MockedClusterMetrics extends ClusterMetricsImpl {
+ public static class MockedClusterMetrics extends ClusterMetricsSnapshot {
/** Some long (> 24h) duration. */
public static final long LONG_DURATION_MS = TimeUnit.DAYS.toMillis(365);
@@ -1964,10 +1963,8 @@ public static class MockedClusterMetrics extends ClusterMetricsImpl {
* @param original - original cluster metrics object. Required to leave the original behaviour for not overriden
* methods.
*/
- public MockedClusterMetrics(ClusterMetricsImpl original) throws Exception {
- super(
- getField(original, "ctx"),
- getField(original, "nodeStartTime"));
+ public MockedClusterMetrics(ClusterMetricsSnapshot original) throws Exception {
+ super(original);
}
/** {@inheritDoc} */