This is an automated email from the ASF dual-hosted git repository.

dockerzhang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/inlong.git


The following commit(s) were added to refs/heads/master by this push:
     new 4a892da914 [INLONG-8177][Sort] Improve jdbc connector object 
calculation and Fix filesystem connector report dirty data metrics error (#8179)
4a892da914 is described below

commit 4a892da914769b8009202d9bedf2f6b0dab07de0
Author: Xin Gong <[email protected]>
AuthorDate: Thu Jun 8 09:56:50 2023 +0800

    [INLONG-8177][Sort] Improve jdbc connector object calculation and Fix 
filesystem connector report dirty data metrics error (#8179)
---
 .../inlong/sort/filesystem/stream/AbstractStreamingWriter.java       | 2 +-
 .../inlong/sort/jdbc/internal/TableMetricStatementExecutor.java      | 5 ++---
 2 files changed, 3 insertions(+), 4 deletions(-)

diff --git 
a/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/filesystem/src/main/java/org/apache/inlong/sort/filesystem/stream/AbstractStreamingWriter.java
 
b/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/filesystem/src/main/java/org/apache/inlong/sort/filesystem/stream/AbstractStreamingWriter.java
index 460d1e84db..293531bdb8 100644
--- 
a/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/filesystem/src/main/java/org/apache/inlong/sort/filesystem/stream/AbstractStreamingWriter.java
+++ 
b/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/filesystem/src/main/java/org/apache/inlong/sort/filesystem/stream/AbstractStreamingWriter.java
@@ -247,7 +247,7 @@ public abstract class AbstractStreamingWriter<IN, OUT> 
extends AbstractStreamOpe
                 throw new RuntimeException(e);
             }
             if (sinkMetricData != null) {
-                sinkMetricData.invokeWithEstimate(element.getValue());
+                sinkMetricData.invokeDirtyWithEstimate(element.getValue());
             }
             if (dirtySink != null) {
                 DirtyData.Builder<Object> builder = DirtyData.builder();
diff --git 
a/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/jdbc/src/main/java/org/apache/inlong/sort/jdbc/internal/TableMetricStatementExecutor.java
 
b/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/jdbc/src/main/java/org/apache/inlong/sort/jdbc/internal/TableMetricStatementExecutor.java
index 8e294314ad..92d572175e 100644
--- 
a/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/jdbc/src/main/java/org/apache/inlong/sort/jdbc/internal/TableMetricStatementExecutor.java
+++ 
b/inlong-sort/sort-flink/sort-flink-v1.13/sort-connectors/jdbc/src/main/java/org/apache/inlong/sort/jdbc/internal/TableMetricStatementExecutor.java
@@ -31,7 +31,6 @@ import org.apache.flink.table.data.RowData;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
-import java.nio.charset.StandardCharsets;
 import java.sql.Connection;
 import java.sql.SQLException;
 import java.util.ArrayList;
@@ -146,7 +145,7 @@ public final class TableMetricStatementExecutor implements 
JdbcBatchStatementExe
         int writtenSize = errorPositions.get(0);
         long writtenBytes = 0L;
         if (writtenSize > 0) {
-            writtenBytes = (long) 
batch.get(0).toString().getBytes(StandardCharsets.UTF_8).length * writtenSize;
+            writtenBytes = CalculateObjectSizeUtils.getDataSize(batch.get(0)) 
* writtenSize;
         }
         if (!multipleSink) {
             sinkMetricData.invoke(writtenSize, writtenBytes);
@@ -186,7 +185,7 @@ public final class TableMetricStatementExecutor implements 
JdbcBatchStatementExe
                     sinkMetricData.invokeWithEstimate(rowData);
                 } else {
                     metric[0] += 1;
-                    metric[1] += rowData.toString().getBytes().length;
+                    metric[1] += CalculateObjectSizeUtils.getDataSize(rowData);
                 }
             } catch (Exception e) {
                 st.clearParameters();

Reply via email to