{"record":{"id":"b270e0a6a50534ef","repo":"apache/seatunnel","slug":"failed-to-consume-python-source-stdout","errorCode":null,"errorMessage":"Failed to consume python source stdout","messagePattern":"Failed to consume python source stdout","errorType":"exception","errorClass":"IllegalStateException","httpStatus":null,"severity":"critical","filePath":"seatunnel-connectors-v2/connector-python/src/main/java/org/apache/seatunnel/connectors/seatunnel/python/source/PythonSourceReader.java","lineNumber":623,"sourceCode":"                \"Timed out waiting for Python source stdout to close after the process exited; ensure child processes do not inherit stdout\");\n    }\n\n    /** Gives a pump a short grace period to consume bytes already written by the direct child. */\n    private void waitForPumpDrain(Thread thread) throws IOException {\n        if (thread == null || !thread.isAlive()) {\n            return;\n        }\n        try {\n            thread.join(PROCESS_EXIT_CHECK_TIMEOUT_MILLIS);\n        } catch (InterruptedException e) {\n            Thread.currentThread().interrupt();\n            throw new IOException(\"Interrupted while draining python source process output\", e);\n        }\n    }\n\n    private void checkPumpFailures() {\n        if (stdoutPumpFailure != null) {\n            throw new IllegalStateException(\n                    \"Failed to consume python source stdout\", stdoutPumpFailure);\n        }\n        if (stderrPumpFailure == null) {\n            return;\n        }\n        throw new IllegalStateException(\n                \"Failed to consume python source stderr\", stderrPumpFailure);\n    }\n\n    /** Registers one poll so close cannot return while rows or completion are still emitted. */\n    private boolean beginPoll() {\n        synchronized (lifecycleLock) {\n            if (noMoreSplits || closeRequested) {\n                return false;\n            }\n            activePolls++;\n            return true;\n        }","sourceCodeStart":605,"sourceCodeEnd":641,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-python/src/main/java/org/apache/seatunnel/connectors/seatunnel/python/source/PythonSourceReader.java#L605-L641","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","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"],"exampleFix":"// before\nthrow new IllegalStateException(\n        \"Failed to consume python source stdout\", stdoutPumpFailure);\n// after\nThrowable cause = stdoutPumpFailure.getCause();\nLOG.error(\"python stdout pump failed; check that the python script writes only rows to stdout\", cause);\nthrow new IllegalStateException(\n        \"Failed to consume python source stdout\", stdoutPumpFailure);","handlingStrategy":"try-catch","validationCode":"// before submitting: verify python script streams valid output\nProcess p = new ProcessBuilder(pythonScriptArgs).start();\ntry (BufferedReader r = new BufferedReader(new InputStreamReader(p.getInputStream()))) {\n    String line;\n    while ((line = r.readLine()) != null) { /* validate row format */ }\n    int code = p.waitFor();\n    if (code != 0) throw new IllegalStateException(\"python script exited \" + code);\n}","typeGuard":null,"tryCatchPattern":"try {\n    reader.pollNext();\n} catch (IllegalStateException e) {\n    if (e.getMessage().contains(\"Failed to consume python source stdout\")) {\n        LOG.error(\"python stdout pump failed\", e.getCause());\n        // fail/restart the split or job\n    } else throw e;\n}","preventionTips":["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"],"tags":["python","subprocess","stream-consumption","process-crash"],"backgroundTag":"broken-pipe","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}