pvillard31 commented on code in PR #11712:
URL: https://github.com/apache/nifi/pull/11712#discussion_r4080899619


##########
nifi-extension-bundles/nifi-aws-bundle/nifi-aws-kinesis/src/main/java/org/apache/nifi/processors/aws/kinesis/ConsumeKinesis.java:
##########
@@ -946,6 +949,39 @@ private WriteResult writeResults(final ProcessSession 
session, final ProcessCont
         return new WriteResult(produced, parseFailures, totalRecordCount, 
totalBytesConsumed, maxMillisBehind);
     }
 
+    private void recordConsumptionMetrics(final ProcessSession session, final 
List<ShardFetchResult> shardResults) {
+        long recordCount = 0;
+        long bytesConsumed = 0;
+        long maxMillisBehind = -1;
+        String shardId = null;
+        for (final ShardFetchResult result : shardResults) {
+            shardId = result.shardId();
+            maxMillisBehind = Math.max(maxMillisBehind, 
result.millisBehindLatest());
+            for (final UserRecord record : result.records()) {
+                recordCount++;
+                bytesConsumed += record.data().length;
+            }
+        }
+
+        if (recordCount == 0) {
+            return;
+        }
+
+        final Map<String, String> attributes = getMetricAttributes(streamName, 
shardId);
+        
session.adjustCounter(KinesisMetricName.RECORDS_CONSUMED.getMetricName(), 
recordCount, attributes, CommitTiming.NOW);
+        
session.adjustCounter(KinesisMetricName.BYTES_CONSUMED.getMetricName(), 
bytesConsumed, attributes, CommitTiming.NOW);
+        // Kinesis uses -1 when millisBehindLatest is absent so record that as 0

Review Comment:
   Should an absent `MillisBehindLatest` value skip the gauge instead of 
converting the `-1` sentinel to zero?



##########
nifi-extension-bundles/nifi-aws-bundle/nifi-aws-kinesis/src/main/java/org/apache/nifi/processors/aws/kinesis/ConsumeKinesis.java:
##########
@@ -946,6 +949,39 @@ private WriteResult writeResults(final ProcessSession 
session, final ProcessCont
         return new WriteResult(produced, parseFailures, totalRecordCount, 
totalBytesConsumed, maxMillisBehind);
     }
 
+    private void recordConsumptionMetrics(final ProcessSession session, final 
List<ShardFetchResult> shardResults) {
+        long recordCount = 0;
+        long bytesConsumed = 0;
+        long maxMillisBehind = -1;
+        String shardId = null;
+        for (final ShardFetchResult result : shardResults) {
+            shardId = result.shardId();
+            maxMillisBehind = Math.max(maxMillisBehind, 
result.millisBehindLatest());
+            for (final UserRecord record : result.records()) {
+                recordCount++;
+                bytesConsumed += record.data().length;
+            }
+        }
+
+        if (recordCount == 0) {

Review Comment:
   Should the lag gauge be recorded independently of `recordCount`, using the 
shard lag observations, so an empty successful poll can update a previously 
positive lag to zero?



##########
nifi-extension-bundles/nifi-aws-bundle/nifi-aws-kinesis/src/test/java/org/apache/nifi/processors/aws/kinesis/ConsumeKinesisTest.java:
##########
@@ -301,8 +340,68 @@ void testMultipleShardsNoDataLoss() throws Exception {
             
shardsSeen.add(flowFile.getAttribute(ConsumeKinesis.ATTR_SHARD_ID));
             totalRecords += 
Long.parseLong(flowFile.getAttribute("record.count"));
         }
-        assertEquals(Set.of("shard-A", "shard-B"), shardsSeen);
+        assertEquals(Set.of(FIRST_SHARD_ID, SECOND_SHARD_ID), shardsSeen);
         assertEquals(4, totalRecords);
+
+        final Map<String, String> firstAttributes = 
getMetricAttributes(FIRST_SHARD_ID);
+        final Map<String, String> secondAttributes = 
getMetricAttributes(SECOND_SHARD_ID);
+        final long firstShardBytes = payloadBytes(FIRST_JSON_RECORD) + 
payloadBytes(SECOND_JSON_RECORD);
+        final long secondShardBytes = payloadBytes(THIRD_JSON_RECORD) + 
payloadBytes(FOURTH_JSON_RECORD);
+
+        assertEquals(2L, 
runner.getCounterValue(KinesisMetricName.RECORDS_CONSUMED.getMetricName(), 
firstAttributes));
+        assertEquals(2L, 
runner.getCounterValue(KinesisMetricName.RECORDS_CONSUMED.getMetricName(), 
secondAttributes));
+        assertEquals(firstShardBytes, 
runner.getCounterValue(KinesisMetricName.BYTES_CONSUMED.getMetricName(), 
firstAttributes));
+        assertEquals(secondShardBytes, 
runner.getCounterValue(KinesisMetricName.BYTES_CONSUMED.getMetricName(), 
secondAttributes));
+        assertEquals(List.of((double) FIRST_SHARD_BEHIND_MS),
+                
runner.getGaugeValues(KinesisMetricName.CONSUMER_MILLISECONDS_BEHIND.getMetricName(),
 firstAttributes));
+        assertEquals(List.of((double) SECOND_SHARD_BEHIND_MS),
+                
runner.getGaugeValues(KinesisMetricName.CONSUMER_MILLISECONDS_BEHIND.getMetricName(),
 secondAttributes));
+    }
+
+    @Test
+    void testEmptyConsumeDoesNotRecordMetrics() throws Exception {

Review Comment:
   Can we add a test where lag changes from positive to zero on an empty 
successful response, to verify that the gauge reports the caught-up state?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to