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]