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

github-merge-queue[bot] pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git


The following commit(s) were added to refs/heads/dev by this push:
     new b03cf1f169 [Fix][Zeta] Fix schema-first multi-table sink metrics 
(#12001)
b03cf1f169 is described below

commit b03cf1f1697ed84d0f57f416673cec40a075ef5a
Author: AOA <[email protected]>
AuthorDate: Sun Aug 30 17:58:27 2026 +0000

    [Fix][Zeta] Fix schema-first multi-table sink metrics (#12001)
---
 .../engine/server/task/flow/SinkFlowLifeCycle.java |  18 ++--
 .../task/flow/SinkFlowLifeCycleMetricsTest.java    | 107 +++++++++++++++++++++
 2 files changed, 119 insertions(+), 6 deletions(-)

diff --git 
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/flow/SinkFlowLifeCycle.java
 
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/flow/SinkFlowLifeCycle.java
index ffb6239108..d723a2b167 100644
--- 
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/flow/SinkFlowLifeCycle.java
+++ 
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/flow/SinkFlowLifeCycle.java
@@ -132,8 +132,8 @@ public class SinkFlowLifeCycle<T, CommitInfoT extends 
Serializable, AggregatedCo
 
     private final EventListener eventListener;
 
-    /** Mapping relationship between upstream TablePath and downstream 
TablePath. */
-    private final Map<TablePath, TablePath> tablesMaps = new HashMap<>();
+    /** Mapping relationship between upstream row table IDs and downstream 
table IDs. */
+    private final Map<String, String> sinkTableMappings = new HashMap<>();
 
     private final MetricsContext metricsContext;
 
@@ -196,8 +196,14 @@ public class SinkFlowLifeCycle<T, CommitInfoT extends 
Serializable, AggregatedCo
         List<TablePath> sinkTables = new ArrayList<>();
         boolean isMulti = sinkAction.getSink() instanceof MultiTableSink;
         if (isMulti) {
-            sinkTables = ((MultiTableSink) 
sinkAction.getSink()).getSinkTables();
-            tablesMaps.putAll(((MultiTableSink) 
sinkAction.getSink()).getSinkTableMapping());
+            MultiTableSink multiTableSink = (MultiTableSink) 
sinkAction.getSink();
+            sinkTables = multiTableSink.getSinkTables();
+            multiTableSink
+                    .getSinkTableMapping()
+                    .forEach(
+                            (sourceTable, sinkTable) ->
+                                    sinkTableMappings.put(
+                                            sourceTable.toString(), 
sinkTable.getFullName()));
         } else {
             Optional<CatalogTable> catalogTable = 
sinkAction.getSink().getWriteCatalogTable();
             if (catalogTable.isPresent()) {
@@ -738,8 +744,8 @@ public class SinkFlowLifeCycle<T, CommitInfoT extends 
Serializable, AggregatedCo
             if (row.getTableId() == null || row.getTableId().isEmpty()) {
                 return row.getTableId();
             }
-            TablePath tablePath = 
tablesMaps.get(TablePath.of(row.getTableId()));
-            return tablePath != null ? tablePath.getFullName() : 
TablePath.DEFAULT.getFullName();
+            return sinkTableMappings.getOrDefault(
+                    row.getTableId(), TablePath.DEFAULT.getFullName());
         }
         Optional<CatalogTable> writeCatalogTable = 
this.sinkAction.getSink().getWriteCatalogTable();
         return writeCatalogTable
diff --git 
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/task/flow/SinkFlowLifeCycleMetricsTest.java
 
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/task/flow/SinkFlowLifeCycleMetricsTest.java
new file mode 100644
index 0000000000..76b6facf05
--- /dev/null
+++ 
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/task/flow/SinkFlowLifeCycleMetricsTest.java
@@ -0,0 +1,107 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.engine.server.task.flow;
+
+import org.apache.seatunnel.api.common.metrics.MetricNames;
+import org.apache.seatunnel.api.sink.SinkWriter;
+import 
org.apache.seatunnel.api.sink.multitablesink.MultiTableAggregatedCommitInfo;
+import org.apache.seatunnel.api.sink.multitablesink.MultiTableCommitInfo;
+import org.apache.seatunnel.api.sink.multitablesink.MultiTableSink;
+import org.apache.seatunnel.api.sink.multitablesink.MultiTableState;
+import org.apache.seatunnel.api.table.catalog.TablePath;
+import org.apache.seatunnel.api.table.type.Record;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.engine.common.utils.concurrent.CompletableFuture;
+import org.apache.seatunnel.engine.core.dag.actions.SinkAction;
+import org.apache.seatunnel.engine.server.execution.TaskGroupLocation;
+import org.apache.seatunnel.engine.server.execution.TaskLocation;
+import org.apache.seatunnel.engine.server.metrics.SeaTunnelMetricsContext;
+import org.apache.seatunnel.engine.server.task.SeaTunnelTask;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
+
+import java.lang.reflect.Field;
+import java.util.Collections;
+
+class SinkFlowLifeCycleMetricsTest {
+
+    private static final TaskLocation TASK_LOCATION =
+            new TaskLocation(new TaskGroupLocation(1L, 1, 1L), 2L, 0);
+
+    @Test
+    void schemaFirstTableIdUsesResolvedSinkTableMetrics() throws Exception {
+        TablePath sourceTable = TablePath.of(null, "PDC_SCHEMA", "CUSTOMER");
+        TablePath sinkTable = TablePath.of("target_db", "CUSTOMER");
+        MultiTableSink sink = Mockito.mock(MultiTableSink.class);
+        
Mockito.when(sink.getSinkTables()).thenReturn(Collections.singletonList(sinkTable));
+        Mockito.when(sink.getSinkTableMapping())
+                .thenReturn(Collections.singletonMap(sourceTable, sinkTable));
+        SinkAction<
+                        SeaTunnelRow,
+                        MultiTableState,
+                        MultiTableCommitInfo,
+                        MultiTableAggregatedCommitInfo>
+                action =
+                        new SinkAction<>(
+                                7L, "sink", sink, Collections.emptySet(), 
Collections.emptySet());
+        SeaTunnelTask task = Mockito.mock(SeaTunnelTask.class);
+        Mockito.when(task.isObservabilityEnabled()).thenReturn(true);
+        SeaTunnelMetricsContext metrics = new SeaTunnelMetricsContext();
+        SinkFlowLifeCycle<
+                        SeaTunnelRow,
+                        MultiTableCommitInfo,
+                        MultiTableAggregatedCommitInfo,
+                        MultiTableState>
+                flow =
+                        new SinkFlowLifeCycle<>(
+                                action,
+                                TASK_LOCATION,
+                                0,
+                                task,
+                                null,
+                                false,
+                                new CompletableFuture<>(),
+                                metrics);
+        setWriter(flow, Mockito.mock(SinkWriter.class));
+        SeaTunnelRow row = new SeaTunnelRow(new Object[] {1});
+        row.setTableId(sourceTable.getFullName());
+
+        flow.received(new Record<>(row));
+
+        Assertions.assertEquals(
+                1L,
+                metrics.counter(MetricNames.SINK_WRITE_COUNT + "#" + 
sinkTable.getFullName())
+                        .getCount());
+        Assertions.assertEquals(
+                0L,
+                metrics.counter(
+                                MetricNames.SINK_WRITE_COUNT
+                                        + "#"
+                                        + TablePath.DEFAULT.getFullName())
+                        .getCount());
+    }
+
+    private static void setWriter(SinkFlowLifeCycle<?, ?, ?, ?> flow, 
SinkWriter<?, ?, ?> writer)
+            throws Exception {
+        Field field = SinkFlowLifeCycle.class.getDeclaredField("writer");
+        field.setAccessible(true);
+        field.set(flow, writer);
+    }
+}

Reply via email to