Hi.All.I used Flink1.16 to consume Kafka data and ...
# troubleshooting
d
Hi.All.I used Flink1.16 to consume Kafka data and found a warning exception in the taskmanager. log file:
Copy code
2023-06-14 15:35:36,897 WARN  org.apache.flink.connector.kafka.source.metrics.KafkaSourceReaderMetrics [] - Error when getting Kafka consumer metric "records-lag" for partition "topic-21". Metric "pendingRecords" may not be reported correctly.
java.lang.IllegalStateException: Cannot find Kafka metric matching current filter.
        at org.apache.flink.connector.kafka.MetricUtil.lambda$getKafkaMetric$1(MetricUtil.java:63) ~[flink-iceberg-sink-1.0.jar:?]
        at java.util.Optional.orElseThrow(Optional.java:290) ~[?:1.8.0_202]
        at org.apache.flink.connector.kafka.MetricUtil.getKafkaMetric(MetricUtil.java:61) ~[flink-iceberg-sink-1.0.jar:?]
        at org.apache.flink.connector.kafka.source.metrics.KafkaSourceReaderMetrics.getRecordsLagMetric(KafkaSourceReaderMetrics.java:308) ~[flink-iceberg-sink-1.0.jar:?]
        at org.apache.flink.connector.kafka.source.metrics.KafkaSourceReaderMetrics.lambda$maybeAddRecordsLagMetric$4(KafkaSourceReaderMetrics.java:231) ~[flink-iceberg-sink-1.0.jar:?]
        at java.util.concurrent.ConcurrentHashMap.computeIfAbsent(ConcurrentHashMap.java:1660) [?:1.8.0_202]
        at org.apache.flink.connector.kafka.source.metrics.KafkaSourceReaderMetrics.maybeAddRecordsLagMetric(KafkaSourceReaderMetrics.java:230) [flink-iceberg-sink-1.0.jar:?]
        at org.apache.flink.connector.kafka.source.reader.KafkaPartitionSplitReader.fetch(KafkaPartitionSplitReader.java:139) [flink-iceberg-sink-1.0.jar:?]
        at org.apache.flink.connector.base.source.reader.fetcher.FetchTask.run(FetchTask.java:58) [flink-iceberg-sink-1.0.jar:?]
        at org.apache.flink.connector.base.source.reader.fetcher.SplitFetcher.runOnce(SplitFetcher.java:142) [flink-iceberg-sink-1.0.jar:?]
        at org.apache.flink.connector.base.source.reader.fetcher.SplitFetcher.run(SplitFetcher.java:105) [flink-iceberg-sink-1.0.jar:?]
        at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) [?:1.8.0_202]
        at java.util.concurrent.FutureTask.run(FutureTask.java:266) [?:1.8.0_202]
        at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149) [?:1.8.0_202]
        at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624) [?:1.8.0_202]
        at java.lang.Thread.run(Thread.java:748) [?:1.8.0_202]
Have you ever meet it?
m
No, but that could also be because this seems to be coming from a Flink Iceberg Sink?
d
No, that's just a feature of our jar package. The warnings thrown are all generated on the source side Kafka and need to be viewed in the taskmanger.log in the taskmanager. Directly viewing the exception information of the task is not visible.
m
d
Yes, that's the problem I've met, and I've seen it too, but it hasn't been fundamentally resolved.
m
That ticket is still open, so it's still a known issue
I'll see if I can find someone for this
d
ok
thx