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