apache/seatunnel · critical · IllegalStateException

Failed to consume python source stdout

Error message

Failed to consume python source stdout

What it means

PythonSourceReader spawns a background thread that continuously consumes the Python subprocess's stdout so the OS pipe never blocks the child. If that pump thread dies with a Throwable, it records the failure in stdoutPumpFailure, and checkPumpFailures rethrows it as an IllegalStateException on the next pollNext or verifyProcessExit call. This means the reader can no longer reliably receive rows from the Python process.

Solutions

  1. Check the wrapped cause (getCause()) to find the real pump-thread failure and fix it
  2. Verify the Python script runs to completion without crashing: run it standalone with the same args and environment
  3. Ensure the Python subprocess writes valid, complete rows to stdout and nothing else pollutes stdout (logging must go to stderr)
  4. Increase memory for the child process if it is being OOM-killed
  5. Check the Python interpreter/dependencies versions match what the connector expects

Example fix

// before
throw new IllegalStateException(
        "Failed to consume python source stdout", stdoutPumpFailure);
// after
Throwable cause = stdoutPumpFailure.getCause();
LOG.error("python stdout pump failed; check that the python script writes only rows to stdout", cause);
throw new IllegalStateException(
        "Failed to consume python source stdout", stdoutPumpFailure);
Defensive patterns

Strategy: try-catch

Validate before calling

// before submitting: verify python script streams valid output
Process p = new ProcessBuilder(pythonScriptArgs).start();
try (BufferedReader r = new BufferedReader(new InputStreamReader(p.getInputStream()))) {
    String line;
    while ((line = r.readLine()) != null) { /* validate row format */ }
    int code = p.waitFor();
    if (code != 0) throw new IllegalStateException("python script exited " + code);
}

Try / catch

try {
    reader.pollNext();
} catch (IllegalStateException e) {
    if (e.getMessage().contains("Failed to consume python source stdout")) {
        LOG.error("python stdout pump failed", e.getCause());
        // fail/restart the split or job
    } else throw e;
}

Prevention

When it happens

Trigger: The stdout pump thread throws while reading or deserializing the Python process output (line ~437 sets stdoutPumpFailure); the failure surfaces on the next call to pollNext or verifyProcessExit via checkPumpFailures.

Common situations: Python script crashes mid-stream leaving a truncated/broken output stream; stdout encoding issues; the pipe is broken because the child was killed (OOM-killed, SIGKILL); deserialization of a row failed in the pump thread.

Related errors


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

Appendix: source

Thrown at seatunnel-connectors-v2/connector-python/src/main/java/org/apache/seatunnel/connectors/seatunnel/python/source/PythonSourceReader.java:623

                "Timed out waiting for Python source stdout to close after the process exited; ensure child processes do not inherit stdout");
    }

    /** Gives a pump a short grace period to consume bytes already written by the direct child. */
    private void waitForPumpDrain(Thread thread) throws IOException {
        if (thread == null || !thread.isAlive()) {
            return;
        }
        try {
            thread.join(PROCESS_EXIT_CHECK_TIMEOUT_MILLIS);
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
            throw new IOException("Interrupted while draining python source process output", e);
        }
    }

    private void checkPumpFailures() {
        if (stdoutPumpFailure != null) {
            throw new IllegalStateException(
                    "Failed to consume python source stdout", stdoutPumpFailure);
        }
        if (stderrPumpFailure == null) {
            return;
        }
        throw new IllegalStateException(
                "Failed to consume python source stderr", stderrPumpFailure);
    }

    /** Registers one poll so close cannot return while rows or completion are still emitted. */
    private boolean beginPoll() {
        synchronized (lifecycleLock) {
            if (noMoreSplits || closeRequested) {
                return false;
            }
            activePolls++;
            return true;
        }

View on GitHub (pinned to cf67b549a7)