This is an automated email from the ASF dual-hosted git repository.

marcuse pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/cassandra.git

commit caa3bd83dcb92a3145b7ef0fd73bbd3708b255bf
Author: Sankalp Kohli <[email protected]>
AuthorDate: Mon Mar 23 10:59:27 2020 +0100

    Expose repair streaming metrics
    
    Patch by Sankalp Kohli; reviewed by Ekaterina Dimitrova for CASSANDRA-15656
---
 CHANGES.txt                                                 |  1 +
 src/java/org/apache/cassandra/metrics/StreamingMetrics.java |  2 ++
 src/java/org/apache/cassandra/streaming/StreamSession.java  | 13 +++++++++++++
 3 files changed, 16 insertions(+)

diff --git a/CHANGES.txt b/CHANGES.txt
index 77d69ca..55243b8 100644
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@ -1,4 +1,5 @@
 4.0-alpha4
+ * Expose repair streaming metrics (CASSANDRA-15656)
  * Set now in seconds in the future for validation repairs (CASSANDRA-15655)
  * Emit metric on preview repair failure (CASSANDRA-15654)
  * Use more appropriate logging levels (CASSANDRA-15661)
diff --git a/src/java/org/apache/cassandra/metrics/StreamingMetrics.java 
b/src/java/org/apache/cassandra/metrics/StreamingMetrics.java
index 793a8c0..80a5e13 100644
--- a/src/java/org/apache/cassandra/metrics/StreamingMetrics.java
+++ b/src/java/org/apache/cassandra/metrics/StreamingMetrics.java
@@ -39,6 +39,8 @@ public class StreamingMetrics
     public static final Counter activeStreamsOutbound = 
Metrics.counter(DefaultNameFactory.createMetricName(TYPE_NAME, 
"ActiveOutboundStreams", null));
     public static final Counter totalIncomingBytes = 
Metrics.counter(DefaultNameFactory.createMetricName(TYPE_NAME, 
"TotalIncomingBytes", null));
     public static final Counter totalOutgoingBytes = 
Metrics.counter(DefaultNameFactory.createMetricName(TYPE_NAME, 
"TotalOutgoingBytes", null));
+    public static final Counter totalOutgoingRepairBytes = 
Metrics.counter(DefaultNameFactory.createMetricName(TYPE_NAME, 
"TotalOutgoingRepairBytes", null));
+    public static final Counter totalOutgoingRepairSSTables = 
Metrics.counter(DefaultNameFactory.createMetricName(TYPE_NAME, 
"TotalOutgoingRepairSSTables", null));
     public final Counter incomingBytes;
     public final Counter outgoingBytes;
 
diff --git a/src/java/org/apache/cassandra/streaming/StreamSession.java 
b/src/java/org/apache/cassandra/streaming/StreamSession.java
index 95d3755..05bb5ff 100644
--- a/src/java/org/apache/cassandra/streaming/StreamSession.java
+++ b/src/java/org/apache/cassandra/streaming/StreamSession.java
@@ -510,8 +510,21 @@ public class StreamSession implements 
IEndpointStateChangeSubscriber
         state(State.PREPARING);
         PrepareSynMessage prepare = new PrepareSynMessage();
         prepare.requests.addAll(requests);
+        long totalBytesToStream = 0;
+        long totalSSTablesStreamed = 0;
         for (StreamTransferTask task : transfers.values())
+        {
+            totalBytesToStream += task.getTotalSize();
+            totalSSTablesStreamed += task.getTotalNumberOfFiles();
             prepare.summaries.add(task.getSummary());
+        }
+
+        if(StreamOperation.REPAIR == getStreamOperation())
+        {
+            StreamingMetrics.totalOutgoingRepairBytes.inc(totalBytesToStream);
+            
StreamingMetrics.totalOutgoingRepairSSTables.inc(totalSSTablesStreamed);
+        }
+
         messageSender.sendMessage(prepare);
     }
 


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to