Github user zentol commented on a diff in the pull request:
https://github.com/apache/flink/pull/4769#discussion_r142684336
--- Diff:
flink-connectors/flink-connector-kafka-base/src/main/java/org/apache/flink/streaming/connectors/kafka/internals/AbstractFetcher.java
---
@@ -543,6 +543,18 @@ private void updateMinPunctuatedWatermark(Watermark
nextWatermark) {
// ------------------------- Metrics ----------------------------------
/**
+ * Register offset metrics.
+ */
+ protected MetricGroup registerOffsetMetrics(MetricGroup metricGroup) {
+ if (useMetrics) {
+ MetricGroup kafkaMetricGroup =
metricGroup.addGroup("KafkaConsumer");
+ addOffsetStateGauge(kafkaMetricGroup);
+ return kafkaMetricGroup;
+ }
+ return null;
--- End diff --
Its better to return a `UnregisteredMetricsGroup` since it cannot cause
exceptions, which would also renver the changes in `KafkaConsumerThread`
unnecessary.
---