This is an automated email from the ASF dual-hosted git repository. dockerzhang pushed a commit to branch release-1.3.0 in repository https://gitbox.apache.org/repos/asf/inlong.git
commit 98a13686348cb2971fdcafccf0e5d0c5e962a96a Author: Charles <[email protected]> AuthorDate: Tue Aug 23 17:36:29 2022 +0800 [INLONG-5637][Sort] Fix kafka load node npe error (#5638) * [INLONG-5637][Sort] Fix kafka load node npe error * [INLONG-5637][Sort] Fix kafka load node npe error --- .../java/org/apache/inlong/sort/kafka/FlinkKafkaProducer.java | 8 ++------ 1 file changed, 2 insertions(+), 6 deletions(-) diff --git a/inlong-sort/sort-connectors/kafka/src/main/java/org/apache/inlong/sort/kafka/FlinkKafkaProducer.java b/inlong-sort/sort-connectors/kafka/src/main/java/org/apache/inlong/sort/kafka/FlinkKafkaProducer.java index 93e7a8b31..75b965340 100644 --- a/inlong-sort/sort-connectors/kafka/src/main/java/org/apache/inlong/sort/kafka/FlinkKafkaProducer.java +++ b/inlong-sort/sort-connectors/kafka/src/main/java/org/apache/inlong/sort/kafka/FlinkKafkaProducer.java @@ -928,19 +928,15 @@ public class FlinkKafkaProducer<IN> } private void sendOutMetrics(Long rowSize, Long dataSize) { - if (metricData.getNumRecordsOut() != null) { + if (metricData != null) { metricData.getNumRecordsOut().inc(rowSize); - } - if (metricData.getNumBytesOut() != null) { metricData.getNumBytesOut().inc(dataSize); } } private void sendDirtyMetrics(Long rowSize, Long dataSize) { - if (metricData.getDirtyRecords() != null) { + if (metricData != null) { metricData.getDirtyRecords().inc(rowSize); - } - if (metricData.getDirtyBytes() != null) { metricData.getDirtyBytes().inc(dataSize); } }
