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