apache/seatunnel · warning

Failed to close flush timer for task

Error message

Failed to close flush timer for task {}

What it means

SourceFlowLifeCycle.closeFlushTimer (invoked from close()) asks the TaskExecutionService to close the periodic flush timer task for this source task. Failures are logged at WARN and non-fatal; flushFuture is nulled regardless. A lingering timer could produce harmless later ticks or minor resource retention.

Solutions

  1. Typically safe to ignore if it appears during job shutdown; verify no timer tasks remain via engine logs/monitoring.
  2. If it occurs mid-run or repeatedly, check TaskExecutionService state and node logs for scheduling failures.
  3. Avoid cancelling jobs twice concurrently; use one graceful stop path.
  4. Upgrade if a timer-vs-teardown race is fixed in your target SeaTunnel version.

Example fix

// before: double close attempts on the timer
sourceFlowLifeCycle.close();
sourceFlowLifeCycle.close(); // second close fails timer cleanup
// after: guard with idempotent close (engine handles via flushFuture = null) and cancel once
seaTunnel.sh -s <jobId> // single graceful stop
Defensive patterns

Strategy: fallback

Try / catch

// non-fatal: engine logs and nulls flushFuture
try {
    job.execute();
} finally {
    // verify no residual timer tasks via engine metrics
}

Prevention

When it happens

Trigger: closeTimerFlushTask(currentTaskLocation) throws, e.g. the timer task was already removed, the task has been cancelled concurrently, or the TaskExecutionService is shutting down.

Common situations: Job cancellation racing with source close; node shutdown while closing tasks; repeated close attempts on the same flow lifecycle.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/435205967b9a1525. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/flow/SourceFlowLifeCycle.java:394

                        .registerTimerFlushTask(
                                currentTaskLocation, this::onTimerTick, flushIntervalMs);
        log.info(
                "Registered flush timer for source task {}, intervalMs={}",
                currentTaskLocation,
                flushIntervalMs);
    }

    private void closeFlushTimer() {
        if (flushFuture == null) {
            return;
        }
        try {
            runningTask
                    .getExecutionContext()
                    .getTaskExecutionService()
                    .closeTimerFlushTask(currentTaskLocation);
        } catch (Exception e) {
            log.warn("Failed to close flush timer for task {}", currentTaskLocation, e);
        }
        flushFuture = null;
    }

    /**
     * Sends a split request to the remote split enumerator.
     *
     * <p>Sends a {@link RequestSplitOperation} to the enumerator, requesting new splits to be
     * assigned to this reader. The enumerator will respond asynchronously by calling {@link
     * #receivedSplits(List)}.
     *
     * @throws RuntimeException if the split request fails due to communication errors
     */
    public void requestSplit() {
        try {
            runningTask
                    .getExecutionContext()
                    .sendToMember(

View on GitHub (pinned to cf67b549a7)