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
- Check the wrapped cause (getCause()) to find the real pump-thread failure and fix it
- Verify the Python script runs to completion without crashing: run it standalone with the same args and environment
- Ensure the Python subprocess writes valid, complete rows to stdout and nothing else pollutes stdout (logging must go to stderr)
- Increase memory for the child process if it is being OOM-killed
- 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
- Ensure the python script writes only well-formed rows to stdout; send logs to stderr
- Run the python script standalone with identical args/env before deploying
- Watch for OS OOM-kills of the child process (check dmesg) and provision enough memory
- Pin python interpreter and dependency versions across all worker nodes
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
- Failed to consume python source stderr
- Timed out waiting for Python source stdout to close after…
- ANTHROPIC_API_KEY environment variable is required for…
- anthropic package required for AI_PROVIDER=anthropic…
- bedrock-mantle provider requires: pip install openai…
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)