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} */