From e6902b00a00a4a1f2b3910c709c0a425d04b69cc Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 24 Jun 2026 14:34:40 +0800 Subject: [PATCH] feat: add spark-submit output parser --- spark_executor/core/log_parser.py | 33 +++++++++++++++++++++++++++++++ tests/unit/test_log_parser.py | 32 ++++++++++++++++++++++++++++++ 2 files changed, 65 insertions(+) create mode 100644 spark_executor/core/log_parser.py create mode 100644 tests/unit/test_log_parser.py diff --git a/spark_executor/core/log_parser.py b/spark_executor/core/log_parser.py new file mode 100644 index 0000000..ad1b72e --- /dev/null +++ b/spark_executor/core/log_parser.py @@ -0,0 +1,33 @@ +# coding=utf-8 +""" +@Time :2026/6/24 +@Author :tao.chen +""" +import re + +_APP_ID_RE = re.compile(r"Submitted application (\S+)") +_TRACKING_URL_RE = re.compile(r"tracking URL:\s+(\S+)") + + +def parse_spark_submit_output(stderr: str) -> tuple[str, str | None]: + """Extract (application_id, tracking_url) from spark-submit stderr. + + Raises ValueError if no application_id can be found. + """ + app_match = _APP_ID_RE.search(stderr) + if app_match: + application_id = app_match.group(1) + else: + url_match = _TRACKING_URL_RE.search(stderr) + if not url_match: + raise ValueError("Could not find application_id in spark-submit output") + # tracking URL is of the form http://rm:8088/proxy/application_xxx/ + tracking_url = url_match.group(1) + tail = tracking_url.rstrip("/").rsplit("/", 1)[-1] + if not tail.startswith("application_"): + raise ValueError("Could not find application_id in spark-submit output") + application_id = tail + + url_match = _TRACKING_URL_RE.search(stderr) + tracking_url = url_match.group(1) if url_match else None + return application_id, tracking_url diff --git a/tests/unit/test_log_parser.py b/tests/unit/test_log_parser.py new file mode 100644 index 0000000..7411fdd --- /dev/null +++ b/tests/unit/test_log_parser.py @@ -0,0 +1,32 @@ +# coding=utf-8 +import pytest + +from spark_executor.core.log_parser import parse_spark_submit_output + + +def test_parses_submitted_application_line(): + stderr = "Warning: ignoring...\nSubmitted application application_17400000001\n" + app_id, url = parse_spark_submit_output(stderr) + assert app_id == "application_17400000001" + assert url is None + + +def test_parses_tracking_url_line(): + stderr = "tracking URL: http://rm:8088/proxy/application_17400000002/\n" + app_id, url = parse_spark_submit_output(stderr) + assert app_id == "application_17400000002" + assert url == "http://rm:8088/proxy/application_17400000002/" + + +def test_prefers_submitted_application_line_over_tracking_url(): + stderr = ( + "tracking URL: http://rm:8088/proxy/application_9999/\n" + "Submitted application application_1234\n" + ) + app_id, _ = parse_spark_submit_output(stderr) + assert app_id == "application_1234" + + +def test_raises_when_no_application_id_found(): + with pytest.raises(ValueError, match="application_id"): + parse_spark_submit_output("some unrelated output\n")