From e4fb97c5cc41a3c7ca704bc705163b690c96dff1 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 30 Jun 2026 14:50:09 +0800 Subject: [PATCH] fix(submit): salvage application_id from spark-submit failure stderr MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Spark-submit can fail AFTER YARN has already accepted the application (lost RM connection, PySpark script exited non-zero on the cluster side, auth expired mid-submit, ...). In those cases the stderr still contains 'Submitted application ', but the previous confirm_submit_job implementation threw away that signal and recorded the pending as FAILED with no application_id. The user had no way to fetch the YARN logs of the failed job. Fix: - run_spark_submit now attaches the failed CompletedProcess to the SparkSubmitError as .result, so callers can still inspect stderr. - confirm_submit_job's failure branch tries to parse (application_id, tracking_url) from exc.result.stderr. If found, it writes them onto the pending record before re-raising. The status is still FAILED and the error message is still set — the user sees the failure normally, but the pending now also has a log target they can pass to get_application_logs. - If the stderr has no application_id (e.g. spark-submit failed locally before reaching YARN), the parse attempt returns None and the pending is FAILED with no log target — exactly the previous behavior in that case. Tests: - test_confirm_marks_failed_on_spark_submit_error: existing test now also asserts application_id is None (no result attached). - test_confirm_recovers_application_id_from_failed_spark_submit: failed run with 'Submitted application X' in stderr → pending is FAILED + application_id is set. - test_confirm_failed_spark_submit_without_application_id_keeps_it_none: failed run with no application_id in stderr → pending is FAILED + application_id is None. Full suite 248 passed (+2 from 246). Co-Authored-By: Claude Fable 5 --- spark_executor/core/spark_submit.py | 14 +++++- spark_executor/tools/submit.py | 24 +++++++++- tests/unit/test_submit_tool.py | 72 +++++++++++++++++++++++++++++ 3 files changed, 108 insertions(+), 2 deletions(-) diff --git a/spark_executor/core/spark_submit.py b/spark_executor/core/spark_submit.py index bffab06..0692bf4 100644 --- a/spark_executor/core/spark_submit.py +++ b/spark_executor/core/spark_submit.py @@ -48,6 +48,13 @@ def build_spark_submit_command( def run_spark_submit(cmd: list[str]) -> "subprocess.CompletedProcess[str]": + """Invoke `spark-submit` and return the CompletedProcess. + + On non-zero exit, raises SparkSubmitError with the failed result + attached as `exc.result` so callers that want to salvage information + from stderr (notably the YARN application_id, which is often present + even when spark-submit itself failed) can do so. + """ logger.debug(f"run_spark_submit exec: {cmd}") result = subprocess.run(cmd, capture_output=True, text=True, errors="replace") logger.debug( @@ -56,7 +63,12 @@ def run_spark_submit(cmd: list[str]) -> "subprocess.CompletedProcess[str]": ) if result.returncode != 0: logger.error(f"spark-submit failed (rc={result.returncode}): {result.stderr[:500]}") - raise SparkSubmitError( + err = SparkSubmitError( f"spark-submit failed (rc={result.returncode}): {result.stderr}" ) + # Attach the failed result so callers can still parse stderr for + # the YARN application_id and surface it in the pending record — + # see confirm_submit_job's failure branch. + err.result = result + raise err return result diff --git a/spark_executor/tools/submit.py b/spark_executor/tools/submit.py index c64150f..f8014ff 100644 --- a/spark_executor/tools/submit.py +++ b/spark_executor/tools/submit.py @@ -208,11 +208,33 @@ def confirm_submit_job(*, pending_id: str) -> SubmitResult: try: result = run_spark_submit(cmd) except SparkSubmitError as exc: + # spark-submit can fail for many reasons (YARN RM unreachable, + # auth expired, the YARN app was launched but spark-submit lost + # its connection, the PySpark script exited non-zero after + # YARN had already accepted it, ...). In several of those cases + # the stderr still contains "Submitted application " — try + # to salvage it so the user can fetch logs via + # get_application_logs on the failed job. Failure to recover + # application_id is non-fatal; we still mark FAILED and raise. + if getattr(exc, "result", None) is not None: + try: + application_id, tracking_url = parse_spark_submit_output( + exc.result.stderr + ) + except (ValueError, AttributeError, TypeError): + application_id, tracking_url = None, None + else: + application_id, tracking_url = None, None + pending.status = "FAILED" pending.error = str(exc) + if application_id: + pending.application_id = application_id + pending.tracking_url = tracking_url pending_store.save(pending) logger.error( - f"confirm_submit_job failed pending_id={pending_id} err={exc}" + f"confirm_submit_job failed pending_id={pending_id} " + f"application_id={application_id} err={exc}" ) raise diff --git a/tests/unit/test_submit_tool.py b/tests/unit/test_submit_tool.py index cb7a645..f1b15c1 100644 --- a/tests/unit/test_submit_tool.py +++ b/tests/unit/test_submit_tool.py @@ -360,6 +360,78 @@ def test_confirm_marks_failed_on_spark_submit_error(monkeypatch, real_script): p = submit.pending_store.get(pid) assert p.status == "FAILED" assert "boom" in (p.error or "") + assert p.application_id is None # no process attached, nothing to recover + + +def test_confirm_recovers_application_id_from_failed_spark_submit(monkeypatch, real_script): + """When spark-submit fails AFTER YARN accepted the application (a + common case: spark-submit loses its connection to the RM, or the + PySpark script exits non-zero on the cluster side), the stderr + still contains 'Submitted application '. We must salvage it so + the user can call get_application_logs on the failed job.""" + submit.prepare_submit_job( + connection="prod", + script_path=str(real_script), + queue="default", + executor_memory="2G", + executor_cores=2, + num_executors=2, + app_name="test-app", + ) + pid = _last_pending_id() + failed_stderr = ( + "Warning: ...\n" + "Submitted application application_17400000002_0007\n" + "tracking URL: http://rm:8088/proxy/application_17400000002_0007/\n" + "... then spark-submit died for some reason\n" + ) + failed_proc = type("P", (), {"returncode": 1, "stderr": failed_stderr})() + def _raise_with_result(_cmd): + err = SparkSubmitError("lost connection to RM") + err.result = failed_proc + raise err + monkeypatch.setattr(submit, "run_spark_submit", _raise_with_result) + with pytest.raises(SparkSubmitError, match="lost connection to RM"): + submit.confirm_submit_job(pending_id=pid) + p = submit.pending_store.get(pid) + assert p.status == "FAILED" + assert "lost connection to RM" in (p.error or "") + # The whole point of this test: application_id must be set even on failure. + assert p.application_id == "application_17400000002_0007" + assert p.tracking_url == "http://rm:8088/proxy/application_17400000002_0007/" + + +def test_confirm_failed_spark_submit_without_application_id_keeps_it_none(monkeypatch, real_script): + """spark-submit can fail BEFORE reaching YARN (e.g. local script + syntax error, missing spark-submit binary). stderr has no + application_id — the pending must record FAILED with no log + target, so get_application_logs will return the clean 'no logs' + error rather than a confusing 404.""" + submit.prepare_submit_job( + connection="prod", + script_path=str(real_script), + queue="default", + executor_memory="2G", + executor_cores=2, + num_executors=2, + app_name="test-app", + ) + pid = _last_pending_id() + failed_proc = type("P", (), { + "returncode": 2, + "stderr": "Error: script_path does not exist: /tmp/whatever.py\n", + })() + def _raise_with_result(_cmd): + err = SparkSubmitError("local failure") + err.result = failed_proc + raise err + monkeypatch.setattr(submit, "run_spark_submit", _raise_with_result) + with pytest.raises(SparkSubmitError, match="local failure"): + submit.confirm_submit_job(pending_id=pid) + p = submit.pending_store.get(pid) + assert p.status == "FAILED" + assert p.application_id is None + assert p.tracking_url is None # --- confirm idempotency / manual retry ---