exceptionfactory commented on code in PR #11712:
URL: https://github.com/apache/nifi/pull/11712#discussion_r4084773479
##########
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:
Thanks, that makes sense. On balance, it seems better to mirror the
indicator that AWS provides, so I will switch this to check for `-1` and not
report if found
--
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]