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

jamesshao pushed a commit to branch upsert
in repository https://gitbox.apache.org/repos/asf/incubator-pinot.git

commit 5f33a29351367b86d0bd29716610286de97fbd3b
Author: james Shao <[email protected]>
AuthorDate: Fri Nov 8 15:01:23 2019 -0800

    add one more lag metrics to kc
    
    Reviewers: tingchen, bzzhang, #streaming_pinot
    
    Reviewed By: bzzhang, #streaming_pinot
    
    Differential Revision: https://code.uberinternal.com/D3574189
---
 .../java/org/apache/pinot/grigio/common/metrics/GrigioGauge.java  | 1 +
 .../apache/pinot/grigio/common/rpcQueue/QueueConsumerRecord.java  | 8 +++++++-
 .../pinot/grigio/common/rpcQueue/QueueConsumerRecordTest.java     | 3 ++-
 .../apache/pinot/grigio/common/rpcQueue/KafkaQueueConsumer.java   | 3 ++-
 .../keyCoordinator/internal/DistributedKeyCoordinatorCore.java    | 2 ++
 5 files changed, 14 insertions(+), 3 deletions(-)

diff --git 
a/pinot-grigio/pinot-grigio-common/src/main/java/org/apache/pinot/grigio/common/metrics/GrigioGauge.java
 
b/pinot-grigio/pinot-grigio-common/src/main/java/org/apache/pinot/grigio/common/metrics/GrigioGauge.java
index 00c5bfc..f9838c1 100644
--- 
a/pinot-grigio/pinot-grigio-common/src/main/java/org/apache/pinot/grigio/common/metrics/GrigioGauge.java
+++ 
b/pinot-grigio/pinot-grigio-common/src/main/java/org/apache/pinot/grigio/common/metrics/GrigioGauge.java
@@ -31,6 +31,7 @@ public enum GrigioGauge implements AbstractMetrics.Gauge {
   VERSION_PRODUCED("versions", MetricsType.KC_ONLY),
   KC_VERSION_CONSUMED("versions", MetricsType.KC_ONLY),
   SERVER_VERSION_CONSUMED("versions", MetricsType.SERVER_ONLY),
+  KC_INPUT_MESSAGE_LAG_MS("milliseconds", MetricsType.KC_ONLY)
   ;
 
   private final String _gaugeName;
diff --git 
a/pinot-grigio/pinot-grigio-common/src/main/java/org/apache/pinot/grigio/common/rpcQueue/QueueConsumerRecord.java
 
b/pinot-grigio/pinot-grigio-common/src/main/java/org/apache/pinot/grigio/common/rpcQueue/QueueConsumerRecord.java
index e3fa109..7e86937 100644
--- 
a/pinot-grigio/pinot-grigio-common/src/main/java/org/apache/pinot/grigio/common/rpcQueue/QueueConsumerRecord.java
+++ 
b/pinot-grigio/pinot-grigio-common/src/main/java/org/apache/pinot/grigio/common/rpcQueue/QueueConsumerRecord.java
@@ -29,13 +29,15 @@ public class QueueConsumerRecord<K, V> {
   private final long _offset;
   private final K _key;
   private final V _record;
+  private final long _timestamp;
 
-  public QueueConsumerRecord(String topic, int partition, long offset, K key, 
V record) {
+  public QueueConsumerRecord(String topic, int partition, long offset, K key, 
V record, long timestamp) {
     this._topic = topic;
     this._partition = partition;
     this._offset = offset;
     this._key = key;
     this._record = record;
+    this._timestamp = timestamp;
   }
 
   public String getTopic() {
@@ -57,4 +59,8 @@ public class QueueConsumerRecord<K, V> {
   public V getRecord() {
     return _record;
   }
+
+  public long getTimestamp() {
+    return _timestamp;
+  }
 }
diff --git 
a/pinot-grigio/pinot-grigio-common/src/test/java/org/apache/pinot/grigio/common/rpcQueue/QueueConsumerRecordTest.java
 
b/pinot-grigio/pinot-grigio-common/src/test/java/org/apache/pinot/grigio/common/rpcQueue/QueueConsumerRecordTest.java
index 3d7f006..b129d08 100644
--- 
a/pinot-grigio/pinot-grigio-common/src/test/java/org/apache/pinot/grigio/common/rpcQueue/QueueConsumerRecordTest.java
+++ 
b/pinot-grigio/pinot-grigio-common/src/test/java/org/apache/pinot/grigio/common/rpcQueue/QueueConsumerRecordTest.java
@@ -25,12 +25,13 @@ public class QueueConsumerRecordTest {
 
   @Test
   public void testGets() {
-    QueueConsumerRecord<String, String> record = new 
QueueConsumerRecord<>("topic", 1, 2, "key", "record");
+    QueueConsumerRecord<String, String> record = new 
QueueConsumerRecord<>("topic", 1, 2, "key", "record", 123);
     Assert.assertEquals(record.getTopic(), "topic");
     Assert.assertEquals(record.getPartition(), 1);
     Assert.assertEquals(record.getOffset(), 2);
     Assert.assertEquals(record.getKey(), "key");
     Assert.assertEquals(record.getRecord(), "record");
+    Assert.assertEquals(record.getTimestamp(), 123);
 
   }
 }
\ No newline at end of file
diff --git 
a/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/rpcQueue/KafkaQueueConsumer.java
 
b/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/rpcQueue/KafkaQueueConsumer.java
index 1391c22..ab00e39 100644
--- 
a/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/rpcQueue/KafkaQueueConsumer.java
+++ 
b/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/rpcQueue/KafkaQueueConsumer.java
@@ -122,7 +122,8 @@ public abstract class KafkaQueueConsumer<K, V> implements 
QueueConsumer<K, V> {
       msgList = new ArrayList<>(records.count());
       for (ConsumerRecord<K, V> record : records) {
         msgList.add(
-            new QueueConsumerRecord<>(record.topic(), record.partition(), 
record.offset(), record.key(), record.value()));
+            new QueueConsumerRecord<>(record.topic(), record.partition(), 
record.offset(), record.key(), record.value(),
+                record.timestamp()));
       }
     }
     
getMetrics().addMeteredGlobalValue(GrigioMeter.MESSAGE_INGEST_COUNT_PER_BATCH, 
msgList.size());
diff --git 
a/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/keyCoordinator/internal/DistributedKeyCoordinatorCore.java
 
b/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/keyCoordinator/internal/DistributedKeyCoordinatorCore.java
index 1cdef42..044673f 100644
--- 
a/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/keyCoordinator/internal/DistributedKeyCoordinatorCore.java
+++ 
b/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/keyCoordinator/internal/DistributedKeyCoordinatorCore.java
@@ -228,6 +228,8 @@ public class DistributedKeyCoordinatorCore {
               LOGGER.warn("exception while trying to put message to queue", e);
             }
           });
+          _metrics.setValueOfGlobalGauge(GrigioGauge.KC_INPUT_MESSAGE_LAG_MS,
+              System.currentTimeMillis() - records.get(records.size() - 
1).getTimestamp());
         }
       } catch (Exception ex) {
         LOGGER.error("encountered exception in consumer ingest loop, will 
retry", ex);


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

Reply via email to