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);
         }
     }

Reply via email to