This is an automated email from the ASF dual-hosted git repository.
yunqing 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 c8864f8f4 [INLONG-6379][Sort] Complement iceberg multiple sink metric
data compute (#6472)
c8864f8f4 is described below
commit c8864f8f43b1b1a42a2b647723eed1a2b460c803
Author: thesumery <[email protected]>
AuthorDate: Wed Nov 9 10:39:06 2022 +0800
[INLONG-6379][Sort] Complement iceberg multiple sink metric data compute
(#6472)
---
.../inlong/sort/iceberg/sink/multiple/IcebergMultipleStreamWriter.java | 3 +++
1 file changed, 3 insertions(+)
diff --git
a/inlong-sort/sort-connectors/iceberg/src/main/java/org/apache/inlong/sort/iceberg/sink/multiple/IcebergMultipleStreamWriter.java
b/inlong-sort/sort-connectors/iceberg/src/main/java/org/apache/inlong/sort/iceberg/sink/multiple/IcebergMultipleStreamWriter.java
index 2c67cbd51..56aa77303 100644
---
a/inlong-sort/sort-connectors/iceberg/src/main/java/org/apache/inlong/sort/iceberg/sink/multiple/IcebergMultipleStreamWriter.java
+++
b/inlong-sort/sort-connectors/iceberg/src/main/java/org/apache/inlong/sort/iceberg/sink/multiple/IcebergMultipleStreamWriter.java
@@ -210,6 +210,9 @@ public class IcebergMultipleStreamWriter extends
IcebergProcessFunction<RecordWi
if (multipleWriters.get(tableId) != null) {
for (RowData data : recordWithSchema.getData()) {
multipleWriters.get(tableId).processElement(data);
+ if (metricData != null) {
+ metricData.invokeWithEstimate(data);
+ }
}
} else {
LOG.error("Unregistered table schema for {}.",
recordWithSchema.getTableId());