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/";