pvillard31 commented on code in PR #11712:
URL: https://github.com/apache/nifi/pull/11712#discussion_r4084156004
##########
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:
My concern is that AWS defines 0 as an observed caught-up state, while -1
means no `MillisBehindLatest` value was provided. Converting unknown to zero
would report that the consumer is caught up without evidence, so I think the
gauge should be skipped until a real value is observed. But, to be fair, this
should be an exceptional situation (I assume) so I don't feel strongly about it
and would be OK with the current approach.
--
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]