DanielLeens opened a new issue, #12231: URL: https://github.com/apache/seatunnel/issues/12231
### Search before asking - [x] Searched existing issues/PRs for task teardown leaks, `SeaTunnelTask.close`, and unchecked lifecycle-close failures. No matching fix was found. #10678 changes classloader disposal, not this lifecycle loop. ### What happened `SeaTunnelTask.close()` catches only `IOException` around each flow lifecycle. An unchecked exception from an earlier lifecycle escapes the stream's `forEach`, so later lifecycles are never closed. A later lifecycle can own executor threads, connections, buffers, and plugin references that then outlive the task. This was reproduced on bigdata3 in a separate bounded JVM using the unchanged production `SeaTunnelTask.java` from `dev` commit `75fd4ed4b2e63579a57b461285997951490c4fe7`. The harness supplies two lifecycles: an earlier failing close and a later close that owns a real executor. It does not replace or copy the production close implementation. ### Source and production trigger - [SeaTunnelTask.close](https://github.com/apache/seatunnel/blob/75fd4ed4b2e63579a57b461285997951490c4fe7/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/SeaTunnelTask.java#L449) iterates `allCycles`, catching only `IOException`. - [SinkFlowLifeCycle.close](https://github.com/apache/seatunnel/blob/75fd4ed4b2e63579a57b461285997951490c4fe7/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/flow/SinkFlowLifeCycle.java#L277) calls `writer.close()`. - [MultiTableSinkWriter.close](https://github.com/apache/seatunnel/blob/75fd4ed4b2e63579a57b461285997951490c4fe7/seatunnel-api/src/main/java/org/apache/seatunnel/api/sink/multitablesink/MultiTableSinkWriter.java#L943) explicitly throws `new RuntimeException(firstE[0])` when teardown records an error. Therefore the unchecked failure used in the harness is compatible with an existing production writer path, not an impossible exception type. - The blocking worker's [fallback close](https://github.com/apache/seatunnel/blob/75fd4ed4b2e63579a57b461285997951490c4fe7/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/TaskExecutionService.java#L1166) also catches only `IOException`. Repeating task close does not guarantee recovery: a repeatable failure in the earlier lifecycle skips the later lifecycle again. The required trigger is a task with multiple initialized lifecycles where an earlier close throws an unchecked exception and a later lifecycle owns resources requiring explicit close. The defect does not require all connectors to fail or every job to leak. ### Reproduction results Environment: Linux, OpenJDK 8u502, `-Xmx128m`, 30-second timeout. The current dev class was compiled against existing SeaTunnel runtime dependencies and placed first on the classpath. Source SHA-256: `6711a1c6c0955a95015b05c3a0c22219a001bc8a9b4879ed19eff975c2b65228`. ```text unchecked=false, attempt=1, laterClosed=true, executorShutdown=true unchecked=false, attempt=2, laterClosed=true, executorShutdown=true unchecked=false, attempt=3, laterClosed=true, executorShutdown=true unchecked=true, attempt=1, laterClosed=false, executorShutdown=false unchecked=true, attempt=2, laterClosed=false, executorShutdown=false unchecked=true, attempt=3, laterClosed=false, executorShutdown=false ``` The checked-exception control proves the second lifecycle is reachable and its cleanup works. Only the exception type changes between the two scenarios. The harness explicitly shuts down its executor in `finally`; no test threads or services are left running. <details> <summary>Complete standalone reproduction (run with the engine-server runtime classpath)</summary> ```java /* * 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. */ import org.apache.seatunnel.api.common.metrics.MetricsContext; import org.apache.seatunnel.engine.common.utils.concurrent.CompletableFuture; import org.apache.seatunnel.engine.core.dag.actions.SourceAction; import org.apache.seatunnel.engine.server.dag.physical.config.SourceConfig; import org.apache.seatunnel.engine.server.execution.ProgressState; import org.apache.seatunnel.engine.server.task.SeaTunnelTask; import org.apache.seatunnel.engine.server.task.flow.FlowLifeCycle; import org.apache.seatunnel.engine.server.task.flow.SourceFlowLifeCycle; import java.io.IOException; import java.util.Arrays; import java.util.Collections; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.logging.Logger; /** * Exercises the unchanged production teardown loop with a failing lifecycle before an owned * executor. Every worker created by the harness is explicitly stopped in finally. */ public final class TaskCloseLeakRepro { public static void main(String[] args) throws Exception { run(false); run(true); } private static void run(boolean unchecked) throws Exception { ExecutorService executor = Executors.newSingleThreadExecutor(); CountDownLatch started = new CountDownLatch(1); CountDownLatch stop = new CountDownLatch(1); AtomicBoolean laterClosed = new AtomicBoolean(); executor.submit(() -> { started.countDown(); try { stop.await(); } catch (InterruptedException expected) { Thread.currentThread().interrupt(); } }); if (!started.await(5, TimeUnit.SECONDS)) throw new AssertionError("Worker did not start"); try { FlowLifeCycle first = new FlowLifeCycle() { @Override public void close() throws IOException { if (unchecked) throw new IllegalStateException("injected sink close failure"); throw new IOException("injected checked close failure"); } }; FlowLifeCycle second = new FlowLifeCycle() { @Override public void close() { laterClosed.set(true); executor.shutdownNow(); } }; Task task = new Task(first, second); for (int attempt = 1; attempt <= 3; attempt++) { try { task.close(); throw new AssertionError("Original close failure was swallowed"); } catch (IOException | IllegalStateException expected) { Logger.getLogger("repro").info( "unchecked=" + unchecked + ", attempt=" + attempt + ", laterClosed=" + laterClosed.get() + ", executorShutdown=" + executor.isShutdown()); } if (laterClosed.get() == unchecked || executor.isShutdown() == unchecked) { throw new AssertionError("Unexpected teardown behavior"); } } } finally { stop.countDown(); executor.shutdownNow(); if (!executor.awaitTermination(5, TimeUnit.SECONDS)) { throw new AssertionError("Harness cleanup failed"); } } } /** * Supplies only the lifecycle list; task.close executes the real SeaTunnelTask method and * AbstractTask cleanup, with no Hazelcast member or connector process. */ private static final class Task extends SeaTunnelTask { private Task(FlowLifeCycle... cycles) { super(1L, null, 0, null, Collections.emptyMap()); allCycles = Arrays.asList(cycles); restoreComplete = new CompletableFuture<>(); } @Override protected SourceFlowLifeCycle<?, ?> createSourceFlowLifeCycle( SourceAction<?, ?, ?> action, SourceConfig config, CompletableFuture<Void> future, MetricsContext metrics) { throw new UnsupportedOperationException("No source initialization in this teardown test"); } @Override public ProgressState call() { throw new UnsupportedOperationException("No task scheduling in this teardown test"); } @Override protected void collect() { throw new UnsupportedOperationException("No record processing in this teardown test"); } } } ``` </details> ### Expected behavior and minimum fix Attempt every initialized lifecycle's cleanup even if an earlier close throws a checked or unchecked exception. Preserve the original failure and collect later failures without silently converting task failure into success. Keep existing teardown order and cover first/middle/last failures, multiple failures, and repeated cleanup attempts. ### Relationship to CI evidence PR #11809's old-head engine test failed in `SinkErrorToMysqlIT.testSinkMaxErrorRatioThreshold` because `HikariPool-3 housekeeper` and `st-multi-table-sink-writer-1` remained alive after the two-minute post-job thread check: [failed job](https://github.com/DanielLeens/seatunnel/actions/runs/34185982566/job/102046465852). This prompted the lifecycle audit. The deterministic reproduction above independently proves the teardown bug. It does **not** establish that this is the exact cause of that particular CI failure; that attribution still requires an end-to-end reproduction of the failing job. Do not weaken the thread-leak assertion or increase its timeout as a substitute for cleanup. ### Version / scope - Version: current dev at `75fd4ed4b2e63579a57b461285997951490c4fe7` - Engine: Zeta - Confirmed: resource-owning later lifecycle skipped after unchecked close failure; executor remains live - Not claimed: measured full-cluster heap growth, OOM, or end-to-end reproduction of the CI job -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
