apache / apache/seatunnel

[Bug][Zeta] Unchecked lifecycle close failure skips remaining task resource cleanup

Open
#12,231 3 comments 0 reactions 0 assignees View on GitHub

Nobody has claimed this yet.

bug help wanted Zeta
Dominant language
Java
Stars
9.7k
Forks
2.4k
Avg merge
3d 9h
Merged PRs (30d)
204

Description

### 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.

Complete standalone reproduction (run with the engine-server runtime classpath)

```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 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");
}
}
}

```

### 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

Contributor guide

No contributing guide indexed for this repository

First steps

  1. Read the whole issue, then the project's contributing guide.
  2. Comment on the issue to say you are picking it up — it saves two people doing the same work.
  3. Fork the repository and make your change on a branch.
  4. Open a pull request that references the issue number.

Research direction

Start with SeaTunnelTask.close in seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/SeaTunnelTask.java and compare its lifecycle handling with the fallback close in TaskExecutionService.java. Review SinkFlowLifeCycle.close and MultiTableSinkWriter.close, then run the standalone reproduction to confirm the unchecked-failure case. Done means every initialized lifecycle is attempted, the original failure is preserved, later failures are collected, and first/middle/last, multiple, and repeated cleanup cases are covered.

Written by the indexing model from the issue text.

Assessment

Tech stack
java
Domain
backend
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Active
Clarity
Clearly specified
Newbie friendliness
68/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.