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]
