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 8a1e1de0de [Fix][Core] Propagate transform lifecycle in Flink and 
Spark starters (#12263)
8a1e1de0de is described below

commit 8a1e1de0de5aa04def69c758a5fdd9a3c9b04c5f
Author: zhiweiniu <[email protected]>
AuthorDate: Sat Sep 12 10:59:10 2026 +0000

    [Fix][Core] Propagate transform lifecycle in Flink and Spark starters 
(#12263)
---
 .../flink/execution/TransformExecuteProcessor.java |  72 ++++++++--
 .../execution/TransformExecuteProcessorTest.java   | 151 ++++++++++++++++++++
 .../spark/execution/TransformExecuteProcessor.java |  56 +++++++-
 .../execution/TransformExecuteProcessorTest.java   | 156 +++++++++++++++++++++
 .../e2e/transform/TestPythonTransformIT.java       |   5 +
 5 files changed, 428 insertions(+), 12 deletions(-)

diff --git 
a/seatunnel-core/seatunnel-flink-starter/seatunnel-flink-starter-common/src/main/java/org/apache/seatunnel/core/starter/flink/execution/TransformExecuteProcessor.java
 
b/seatunnel-core/seatunnel-flink-starter/seatunnel-flink-starter-common/src/main/java/org/apache/seatunnel/core/starter/flink/execution/TransformExecuteProcessor.java
index 1fd059022e..6d53b9bc64 100644
--- 
a/seatunnel-core/seatunnel-flink-starter/seatunnel-flink-starter-common/src/main/java/org/apache/seatunnel/core/starter/flink/execution/TransformExecuteProcessor.java
+++ 
b/seatunnel-core/seatunnel-flink-starter/seatunnel-flink-starter-common/src/main/java/org/apache/seatunnel/core/starter/flink/execution/TransformExecuteProcessor.java
@@ -37,12 +37,17 @@ import 
org.apache.seatunnel.plugin.discovery.seatunnel.SeaTunnelFactoryDiscovery
 import 
org.apache.seatunnel.plugin.discovery.seatunnel.SeaTunnelTransformPluginDiscovery;
 
 import org.apache.commons.collections.CollectionUtils;
+import org.apache.flink.api.common.functions.AbstractRichFunction;
 import org.apache.flink.api.common.functions.FlatMapFunction;
+import org.apache.flink.api.common.functions.MapFunction;
 import org.apache.flink.api.common.typeinfo.TypeInformation;
+import org.apache.flink.configuration.Configuration;
 import org.apache.flink.streaming.api.datastream.DataStream;
 import org.apache.flink.streaming.api.operators.StreamMap;
 import org.apache.flink.util.Collector;
 
+import lombok.extern.slf4j.Slf4j;
+
 import java.net.URL;
 import java.util.ArrayList;
 import java.util.LinkedHashMap;
@@ -50,12 +55,14 @@ import java.util.List;
 import java.util.Map;
 import java.util.Objects;
 import java.util.Optional;
+import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.stream.Collectors;
 
 import static 
org.apache.seatunnel.api.options.ConnectorCommonOptions.PLUGIN_NAME;
 import static 
org.apache.seatunnel.api.options.ConnectorCommonOptions.PLUGIN_OUTPUT;
 
 @SuppressWarnings("unchecked,rawtypes")
+@Slf4j
 public class TransformExecuteProcessor
         extends FlinkAbstractPluginExecuteProcessor<TableTransformFactory> {
 
@@ -167,23 +174,72 @@ public class TransformExecuteProcessor
                         new StreamMap<>(
                                 flinkRuntimeEnvironment
                                         .getStreamExecutionEnvironment()
-                                        .clean(
-                                                row ->
-                                                        
((SeaTunnelMapTransform<SeaTunnelRow>)
-                                                                        
transform)
-                                                                .map(row))))
+                                        .clean(new 
TransformMapFunction(transform))))
                 // null value shouldn't be passed to downstream
                 .filter(Objects::nonNull);
     }
 
-    public static class ArrayFlatMap implements FlatMapFunction<SeaTunnelRow, 
SeaTunnelRow> {
+    /**
+     * Bridges Flink's rich-function lifecycle to a {@link SeaTunnelTransform} 
instance.
+     *
+     * <p>Opening and record processing run on the operator task thread, so 
{@code opened} does not
+     * require synchronization. The atomic close guard also makes cleanup safe 
when Flink invokes
+     * more than one cleanup path.
+     */
+    private abstract static class LifecycleAwareTransformFunction extends 
AbstractRichFunction {
 
-        private SeaTunnelTransform transform;
+        protected final SeaTunnelTransform<SeaTunnelRow> transform;
+        private final AtomicBoolean closed = new AtomicBoolean();
+        private transient boolean opened;
 
-        public ArrayFlatMap(SeaTunnelTransform transform) {
+        private 
LifecycleAwareTransformFunction(SeaTunnelTransform<SeaTunnelRow> transform) {
             this.transform = transform;
         }
 
+        @Override
+        public void open(Configuration parameters) {
+            if (!opened) {
+                transform.open();
+                opened = true;
+            }
+        }
+
+        @Override
+        public void close() {
+            // Close at most once and never let a cleanup failure mask the 
task's original outcome.
+            if (!closed.compareAndSet(false, true)) {
+                return;
+            }
+            try {
+                transform.close();
+            } catch (RuntimeException e) {
+                log.warn("Failed to close SeaTunnel transform {}", 
transform.getPluginName(), e);
+            }
+        }
+    }
+
+    /** Lifecycle-aware adapter for map transforms. */
+    static class TransformMapFunction extends LifecycleAwareTransformFunction
+            implements MapFunction<SeaTunnelRow, SeaTunnelRow> {
+
+        TransformMapFunction(SeaTunnelTransform<SeaTunnelRow> transform) {
+            super(transform);
+        }
+
+        @Override
+        public SeaTunnelRow map(SeaTunnelRow row) {
+            return ((SeaTunnelMapTransform<SeaTunnelRow>) transform).map(row);
+        }
+    }
+
+    /** Lifecycle-aware adapter for flat-map transforms. */
+    static class ArrayFlatMap extends LifecycleAwareTransformFunction
+            implements FlatMapFunction<SeaTunnelRow, SeaTunnelRow> {
+
+        ArrayFlatMap(SeaTunnelTransform<SeaTunnelRow> transform) {
+            super(transform);
+        }
+
         @Override
         public void flatMap(SeaTunnelRow row, Collector<SeaTunnelRow> 
collector) {
             List<SeaTunnelRow> rows =
diff --git 
a/seatunnel-core/seatunnel-flink-starter/seatunnel-flink-starter-common/src/test/java/org/apache/seatunnel/core/starter/flink/execution/TransformExecuteProcessorTest.java
 
b/seatunnel-core/seatunnel-flink-starter/seatunnel-flink-starter-common/src/test/java/org/apache/seatunnel/core/starter/flink/execution/TransformExecuteProcessorTest.java
new file mode 100644
index 0000000000..9689dee374
--- /dev/null
+++ 
b/seatunnel-core/seatunnel-flink-starter/seatunnel-flink-starter-common/src/test/java/org/apache/seatunnel/core/starter/flink/execution/TransformExecuteProcessorTest.java
@@ -0,0 +1,151 @@
+/*
+ * 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.core.starter.flink.execution;
+
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.api.transform.SeaTunnelFlatMapTransform;
+import org.apache.seatunnel.api.transform.SeaTunnelMapTransform;
+
+import org.apache.flink.api.common.functions.util.ListCollector;
+import org.apache.flink.configuration.Configuration;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+
+class TransformExecuteProcessorTest {
+
+    @Test
+    void mapFunctionPropagatesLifecycleExactlyOnce() {
+        LifecycleTracker tracker = new LifecycleTracker();
+        TransformExecuteProcessor.TransformMapFunction function =
+                new TransformExecuteProcessor.TransformMapFunction(
+                        new TrackingMapTransform(tracker));
+        SeaTunnelRow row = new SeaTunnelRow(new Object[] {1});
+
+        function.open(new Configuration());
+        function.open(new Configuration());
+        Assertions.assertSame(row, function.map(row));
+        function.close();
+        function.close();
+
+        Assertions.assertEquals(1, tracker.openCount);
+        Assertions.assertEquals(1, tracker.closeCount);
+    }
+
+    @Test
+    void flatMapFunctionPropagatesLifecycleExactlyOnce() throws Exception {
+        LifecycleTracker tracker = new LifecycleTracker();
+        TransformExecuteProcessor.ArrayFlatMap function =
+                new TransformExecuteProcessor.ArrayFlatMap(new 
TrackingFlatMapTransform(tracker));
+        SeaTunnelRow row = new SeaTunnelRow(new Object[] {1});
+        List<SeaTunnelRow> output = new ArrayList<>();
+
+        function.open(new Configuration());
+        function.open(new Configuration());
+        function.flatMap(row, new ListCollector<>(output));
+        function.close();
+        function.close();
+
+        Assertions.assertEquals(Collections.singletonList(row), output);
+        Assertions.assertEquals(1, tracker.openCount);
+        Assertions.assertEquals(1, tracker.closeCount);
+    }
+
+    @Test
+    void closeFailureDoesNotEscapeLifecycleCallback() {
+        LifecycleTracker tracker = new LifecycleTracker();
+        tracker.failOnClose = true;
+        TransformExecuteProcessor.TransformMapFunction function =
+                new TransformExecuteProcessor.TransformMapFunction(
+                        new TrackingMapTransform(tracker));
+
+        function.open(new Configuration());
+        Assertions.assertDoesNotThrow(function::close);
+        Assertions.assertDoesNotThrow(function::close);
+
+        Assertions.assertEquals(1, tracker.closeCount);
+    }
+
+    private static class LifecycleTracker {
+        private int openCount;
+        private int closeCount;
+        private boolean failOnClose;
+    }
+
+    private abstract static class TrackingTransform {
+        protected final LifecycleTracker tracker;
+
+        private TrackingTransform(LifecycleTracker tracker) {
+            this.tracker = tracker;
+        }
+
+        public String getPluginName() {
+            return "lifecycle-test";
+        }
+
+        public void open() {
+            tracker.openCount++;
+        }
+
+        public CatalogTable getProducedCatalogTable() {
+            return null;
+        }
+
+        public List<CatalogTable> getProducedCatalogTables() {
+            return Collections.emptyList();
+        }
+
+        public void close() {
+            tracker.closeCount++;
+            if (tracker.failOnClose) {
+                throw new RuntimeException("expected close failure");
+            }
+        }
+    }
+
+    private static class TrackingMapTransform extends TrackingTransform
+            implements SeaTunnelMapTransform<SeaTunnelRow> {
+
+        private TrackingMapTransform(LifecycleTracker tracker) {
+            super(tracker);
+        }
+
+        @Override
+        public SeaTunnelRow map(SeaTunnelRow row) {
+            return row;
+        }
+    }
+
+    private static class TrackingFlatMapTransform extends TrackingTransform
+            implements SeaTunnelFlatMapTransform<SeaTunnelRow> {
+
+        private TrackingFlatMapTransform(LifecycleTracker tracker) {
+            super(tracker);
+        }
+
+        @Override
+        public List<SeaTunnelRow> flatMap(SeaTunnelRow row) {
+            return Collections.singletonList(row);
+        }
+    }
+}
diff --git 
a/seatunnel-core/seatunnel-spark-starter/seatunnel-spark-starter-common/src/main/java/org/apache/seatunnel/core/starter/spark/execution/TransformExecuteProcessor.java
 
b/seatunnel-core/seatunnel-spark-starter/seatunnel-spark-starter-common/src/main/java/org/apache/seatunnel/core/starter/spark/execution/TransformExecuteProcessor.java
index ee13f1099c..816f4b0668 100644
--- 
a/seatunnel-core/seatunnel-spark-starter/seatunnel-spark-starter-common/src/main/java/org/apache/seatunnel/core/starter/spark/execution/TransformExecuteProcessor.java
+++ 
b/seatunnel-core/seatunnel-spark-starter/seatunnel-spark-starter-common/src/main/java/org/apache/seatunnel/core/starter/spark/execution/TransformExecuteProcessor.java
@@ -40,12 +40,14 @@ import 
org.apache.seatunnel.translation.spark.execution.DatasetTableInfo;
 import org.apache.seatunnel.translation.spark.execution.MultiTableManager;
 
 import org.apache.commons.collections.CollectionUtils;
+import org.apache.spark.TaskContext;
 import org.apache.spark.api.java.function.FlatMapFunction;
 import org.apache.spark.sql.Dataset;
 import org.apache.spark.sql.Row;
 import org.apache.spark.sql.catalyst.encoders.ExpressionEncoder;
 import org.apache.spark.sql.catalyst.encoders.RowEncoder;
 import org.apache.spark.sql.catalyst.expressions.GenericRow;
+import org.apache.spark.util.TaskCompletionListener;
 
 import lombok.extern.slf4j.Slf4j;
 
@@ -57,6 +59,7 @@ import java.util.List;
 import java.util.Map;
 import java.util.Objects;
 import java.util.Optional;
+import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.stream.Collectors;
 
 import static 
org.apache.seatunnel.api.options.ConnectorCommonOptions.PLUGIN_NAME;
@@ -184,10 +187,20 @@ public class TransformExecuteProcessor
                 .filter(Objects::nonNull);
     }
 
-    private static class TransformMapPartitionsFunction implements 
FlatMapFunction<Row, Row> {
-        private SeaTunnelTransform<SeaTunnelRow> transform;
-        private MultiTableManager inputManager;
-        private MultiTableManager outputManager;
+    /**
+     * Applies a transform and connects its lifecycle to the owning Spark task.
+     *
+     * <p>Spark invokes {@link #call(Row)} on the task thread, so the open and 
listener-registration
+     * flags do not require synchronization. The atomic close guard tolerates 
repeated or concurrent
+     * task-completion callbacks.
+     */
+    static class TransformMapPartitionsFunction implements 
FlatMapFunction<Row, Row> {
+        private final SeaTunnelTransform<SeaTunnelRow> transform;
+        private final MultiTableManager inputManager;
+        private final MultiTableManager outputManager;
+        private final AtomicBoolean closed = new AtomicBoolean();
+        private transient boolean completionListenerRegistered;
+        private transient boolean opened;
 
         public TransformMapPartitionsFunction(
                 SeaTunnelTransform<SeaTunnelRow> transform,
@@ -200,6 +213,9 @@ public class TransformExecuteProcessor
 
         @Override
         public Iterator<Row> call(Row row) throws Exception {
+            if (!opened) {
+                initialize(TaskContext.get());
+            }
             List<Row> rows = new ArrayList<>();
 
             SeaTunnelRow seaTunnelRow = inputManager.reconvert((GenericRow) 
row);
@@ -220,5 +236,37 @@ public class TransformExecuteProcessor
             }
             return rows.iterator();
         }
+
+        /**
+         * Opens the transform once for this task.
+         *
+         * <p>The completion listener must be registered before {@link 
SeaTunnelTransform#open()} so
+         * resources created by a partially failed open attempt are still 
released.
+         */
+        void initialize(TaskContext taskContext) {
+            if (opened) {
+                return;
+            }
+            Objects.requireNonNull(taskContext, "Spark TaskContext must be 
available");
+            if (!completionListenerRegistered) {
+                taskContext.addTaskCompletionListener(
+                        (TaskCompletionListener) context -> closeTransform());
+                completionListenerRegistered = true;
+            }
+            transform.open();
+            opened = true;
+        }
+
+        /** Closes the transform at most once without changing the task's 
existing outcome. */
+        private void closeTransform() {
+            if (!closed.compareAndSet(false, true)) {
+                return;
+            }
+            try {
+                transform.close();
+            } catch (RuntimeException e) {
+                log.warn("Failed to close SeaTunnel transform {}", 
transform.getPluginName(), e);
+            }
+        }
     }
 }
diff --git 
a/seatunnel-core/seatunnel-spark-starter/seatunnel-spark-starter-common/src/test/java/org/apache/seatunnel/core/starter/spark/execution/TransformExecuteProcessorTest.java
 
b/seatunnel-core/seatunnel-spark-starter/seatunnel-spark-starter-common/src/test/java/org/apache/seatunnel/core/starter/spark/execution/TransformExecuteProcessorTest.java
new file mode 100644
index 0000000000..df95daacd2
--- /dev/null
+++ 
b/seatunnel-core/seatunnel-spark-starter/seatunnel-spark-starter-common/src/test/java/org/apache/seatunnel/core/starter/spark/execution/TransformExecuteProcessorTest.java
@@ -0,0 +1,156 @@
+/*
+ * 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.core.starter.spark.execution;
+
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.api.transform.SeaTunnelMapTransform;
+
+import org.apache.spark.TaskContext;
+import org.apache.spark.util.TaskCompletionListener;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
+
+import java.util.Collections;
+import java.util.List;
+
+import static org.mockito.Mockito.doReturn;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+
+class TransformExecuteProcessorTest {
+
+    @Test
+    void opensOnceAndClosesOnceOnTaskCompletion() {
+        LifecycleTracker tracker = new LifecycleTracker();
+        TransformExecuteProcessor.TransformMapPartitionsFunction function =
+                createFunction(new TrackingMapTransform(tracker));
+        TaskContext taskContext = mock(TaskContext.class);
+        ArgumentCaptor<TaskCompletionListener> listener = 
registerCompletionListener(taskContext);
+
+        function.initialize(taskContext);
+        function.initialize(taskContext);
+
+        Assertions.assertEquals(1, tracker.openCount);
+        verify(taskContext, 
times(1)).addTaskCompletionListener(listener.capture());
+
+        listener.getValue().onTaskCompletion(taskContext);
+        listener.getValue().onTaskCompletion(taskContext);
+        Assertions.assertEquals(1, tracker.closeCount);
+    }
+
+    @Test
+    void completionListenerSuppressesCloseFailure() {
+        LifecycleTracker tracker = new LifecycleTracker();
+        tracker.failOnClose = true;
+        TransformExecuteProcessor.TransformMapPartitionsFunction function =
+                createFunction(new TrackingMapTransform(tracker));
+        TaskContext taskContext = mock(TaskContext.class);
+        ArgumentCaptor<TaskCompletionListener> listener = 
registerCompletionListener(taskContext);
+
+        function.initialize(taskContext);
+        verify(taskContext).addTaskCompletionListener(listener.capture());
+
+        Assertions.assertDoesNotThrow(() -> 
listener.getValue().onTaskCompletion(taskContext));
+        Assertions.assertEquals(1, tracker.closeCount);
+    }
+
+    @Test
+    void registersCleanupBeforeOpeningTransform() {
+        LifecycleTracker tracker = new LifecycleTracker();
+        tracker.failOnOpen = true;
+        TransformExecuteProcessor.TransformMapPartitionsFunction function =
+                createFunction(new TrackingMapTransform(tracker));
+        TaskContext taskContext = mock(TaskContext.class);
+        ArgumentCaptor<TaskCompletionListener> listener = 
registerCompletionListener(taskContext);
+
+        Assertions.assertThrows(RuntimeException.class, () -> 
function.initialize(taskContext));
+        verify(taskContext).addTaskCompletionListener(listener.capture());
+
+        listener.getValue().onTaskCompletion(taskContext);
+        Assertions.assertEquals(1, tracker.openCount);
+        Assertions.assertEquals(1, tracker.closeCount);
+    }
+
+    private static TransformExecuteProcessor.TransformMapPartitionsFunction 
createFunction(
+            SeaTunnelMapTransform<SeaTunnelRow> transform) {
+        return new 
TransformExecuteProcessor.TransformMapPartitionsFunction(transform, null, null);
+    }
+
+    private static ArgumentCaptor<TaskCompletionListener> 
registerCompletionListener(
+            TaskContext taskContext) {
+        ArgumentCaptor<TaskCompletionListener> listener =
+                ArgumentCaptor.forClass(TaskCompletionListener.class);
+        
doReturn(taskContext).when(taskContext).addTaskCompletionListener(listener.capture());
+        return listener;
+    }
+
+    private static class LifecycleTracker {
+        private int openCount;
+        private int closeCount;
+        private boolean failOnOpen;
+        private boolean failOnClose;
+    }
+
+    private static class TrackingMapTransform implements 
SeaTunnelMapTransform<SeaTunnelRow> {
+        private final LifecycleTracker tracker;
+
+        private TrackingMapTransform(LifecycleTracker tracker) {
+            this.tracker = tracker;
+        }
+
+        @Override
+        public String getPluginName() {
+            return "lifecycle-test";
+        }
+
+        @Override
+        public void open() {
+            tracker.openCount++;
+            if (tracker.failOnOpen) {
+                throw new RuntimeException("expected open failure");
+            }
+        }
+
+        @Override
+        public SeaTunnelRow map(SeaTunnelRow row) {
+            return row;
+        }
+
+        @Override
+        public CatalogTable getProducedCatalogTable() {
+            return null;
+        }
+
+        @Override
+        public List<CatalogTable> getProducedCatalogTables() {
+            return Collections.emptyList();
+        }
+
+        @Override
+        public void close() {
+            tracker.closeCount++;
+            if (tracker.failOnClose) {
+                throw new RuntimeException("expected close failure");
+            }
+        }
+    }
+}
diff --git 
a/seatunnel-e2e/seatunnel-transforms-v2-e2e/seatunnel-transforms-v2-e2e-part-2/src/test/java/org/apache/seatunnel/e2e/transform/TestPythonTransformIT.java
 
b/seatunnel-e2e/seatunnel-transforms-v2-e2e/seatunnel-transforms-v2-e2e-part-2/src/test/java/org/apache/seatunnel/e2e/transform/TestPythonTransformIT.java
index 0591642190..0cb018f230 100644
--- 
a/seatunnel-e2e/seatunnel-transforms-v2-e2e/seatunnel-transforms-v2-e2e-part-2/src/test/java/org/apache/seatunnel/e2e/transform/TestPythonTransformIT.java
+++ 
b/seatunnel-e2e/seatunnel-transforms-v2-e2e/seatunnel-transforms-v2-e2e-part-2/src/test/java/org/apache/seatunnel/e2e/transform/TestPythonTransformIT.java
@@ -31,6 +31,7 @@ import 
org.apache.seatunnel.e2e.common.container.flink.Flink18Container;
 import org.apache.seatunnel.e2e.common.container.flink.Flink20Container;
 import org.apache.seatunnel.e2e.common.container.seatunnel.SeaTunnelContainer;
 import org.apache.seatunnel.e2e.common.junit.ContainerTestingExtension;
+import org.apache.seatunnel.e2e.common.junit.DisabledOnContainer;
 import org.apache.seatunnel.e2e.common.junit.TestCaseInvocationContextProvider;
 import org.apache.seatunnel.e2e.common.junit.TestContainerExtension;
 import org.apache.seatunnel.e2e.common.junit.TestContainers;
@@ -58,6 +59,10 @@ import java.util.stream.Collectors;
     TimingExtension.class
 })
 @TestInstance(TestInstance.Lifecycle.PER_CLASS)
+@DisabledOnContainer(
+        value = {TestContainerId.FLINK_1_13},
+        disabledReason =
+                "The Flink 1.13 image cannot install Python because its Debian 
security metadata is expired")
 public class TestPythonTransformIT {
 
     private static final String BASE_PATH = "/python_transform/";

Reply via email to