Before: SPARK_EXECUTOR_DATA_DIR, SPARK_EXECUTOR_JOBS_DIR, and
YARN_RESOURCE_MANAGER_URL were each read directly via os.environ.get()
inside the module that used them. Log level was hardcoded DEBUG in
common/logging.py. There was no single file showing what the full set
of env vars the app reads is.
After: common/config.py defines a single Settings dataclass that
reads all env vars at import time and exposes them as fields on a
module-level singleton. App code uses "from common.config import
settings; settings.data_dir" etc. New SPARK_EXECUTOR_LOG_LEVEL env
var controls stderr + info file verbosity (debug file always gets full
DEBUG).
Improvements:
- One file lists every env var the app reads (was: grep the codebase)
- Tests can monkeypatch fields on the settings singleton directly
instead of monkeypatching the env + reloading
- Adding a new env var means adding one field in config.py, not
editing 3+ call sites
- settings.reload() method for tests that prefer env-var style
Out of scope (kept where they are):
- GUNICORN_* env vars live in gunicorn.conf.py (gunicorn concept)
- PYTHONUNBUFFERED in Dockerfile (Python runtime flag)
- SPARK_SUBMIT_OPTS not in config (JVM flag, not Python)
Test changes:
- test_job_writer.py: settings.jobs_dir instead of monkeypatching
SPARK_EXECUTOR_JOBS_DIR
- test_yarn_client.py: settings.yarn_resource_manager_url instead of
monkeypatching YARN_RESOURCE_MANAGER_URL
- test_generate_tool.py: same as job_writer
- Each test file gets an autouse fixture that snapshots+restores
settings so one test mutation does not leak into the next
116/116 still pass. Live verified: SPARK_EXECUTOR_LOG_LEVEL=INFO
suppresses DEBUG loguru output as expected.
111 lines
4.5 KiB
Python
111 lines
4.5 KiB
Python
# coding=utf-8
|
|
"""
|
|
@Time :2026/6/24
|
|
@Author :tao.chen
|
|
|
|
YARN ResourceManager REST API client. Replaces the previous `yarn` CLI shell-out
|
|
so the runtime image does not need a Hadoop client installation — the
|
|
`httpx` library already in pyproject.toml is enough.
|
|
|
|
Endpoints used (YARN 2.6+):
|
|
GET /ws/v1/cluster/apps/{appid} -> app status + state
|
|
GET /ws/v1/cluster/apps/{appid}/aggregated-logs -> aggregated container logs
|
|
PUT /ws/v1/cluster/apps/{appid}/state -> kill an app (body: {"state":"KILLED"})
|
|
|
|
The ResourceManager URL is passed in per call (snapshotted on the Job at
|
|
confirm_submit_job time) and falls back to the YARN_RESOURCE_MANAGER_URL env
|
|
var if unset. This matches the pattern the original Connection.yarn_rm_url
|
|
field was designed for, but no longer requires the `yarn` CLI to interpret it.
|
|
"""
|
|
import json
|
|
|
|
import httpx
|
|
|
|
from common.config import settings
|
|
from common.logging import logger
|
|
|
|
|
|
class YarnError(Exception):
|
|
"""Raised when a YARN REST API call fails."""
|
|
|
|
|
|
class YarnConfigError(YarnError):
|
|
"""Raised when the YARN ResourceManager URL is missing or malformed."""
|
|
|
|
|
|
def _base_url(yarn_rm_url: str | None) -> str:
|
|
"""Resolve and validate the RM URL. Raises YarnConfigError if unusable."""
|
|
url = yarn_rm_url or settings.yarn_resource_manager_url
|
|
if not url:
|
|
raise YarnConfigError(
|
|
"YARN ResourceManager URL is not configured. "
|
|
"Set Connection.yarn_rm_url when saving the connection, "
|
|
"or set the YARN_RESOURCE_MANAGER_URL environment variable."
|
|
)
|
|
base = url.rstrip("/")
|
|
if not base.startswith(("http://", "https://")):
|
|
raise YarnConfigError(
|
|
f"YARN ResourceManager URL must start with http:// or https://: {url!r}"
|
|
)
|
|
return base
|
|
|
|
|
|
def _request(method: str, url: str, *, json_body: dict | None = None,
|
|
timeout: float = 30.0) -> httpx.Response:
|
|
logger.debug(f"YARN {method} {url}" + (f" body={json_body}" if json_body else ""))
|
|
try:
|
|
resp = httpx.request(method, url, json=json_body, timeout=timeout)
|
|
except httpx.HTTPError as exc:
|
|
logger.error(f"YARN {method} {url} failed: {exc}")
|
|
raise YarnError(f"YARN connection failed: {exc}") from exc
|
|
logger.debug(
|
|
f"YARN {method} {url} -> {resp.status_code} "
|
|
f"({len(resp.content)} bytes)"
|
|
)
|
|
return resp
|
|
|
|
|
|
def get_application_status(application_id: str, yarn_rm_url: str | None) -> tuple[str, str]:
|
|
"""Return (state, raw_json_text) for an application, or raise YarnError."""
|
|
url = f"{_base_url(yarn_rm_url)}/ws/v1/cluster/apps/{application_id}"
|
|
resp = _request("GET", url)
|
|
if resp.status_code == 404:
|
|
raise YarnError(f"YARN application {application_id!r} not found")
|
|
if resp.status_code >= 400:
|
|
logger.error(f"YARN GET {url} -> {resp.status_code}: {resp.text[:500]}")
|
|
raise YarnError(f"YARN GET returned HTTP {resp.status_code}")
|
|
data = resp.json()
|
|
app = data.get("app", {})
|
|
state = app.get("state")
|
|
if not state:
|
|
raise YarnError(f"Could not parse YARN state from response: {data!r}")
|
|
logger.info(f"YARN status {application_id} -> {state}")
|
|
return state, json.dumps(data, indent=2)
|
|
|
|
|
|
def get_application_logs(application_id: str, yarn_rm_url: str | None) -> str:
|
|
"""Return aggregated container logs for an application as text."""
|
|
url = f"{_base_url(yarn_rm_url)}/ws/v1/cluster/apps/{application_id}/aggregated-logs"
|
|
resp = _request("GET", url, timeout=60.0)
|
|
if resp.status_code == 404:
|
|
raise YarnError(
|
|
f"YARN aggregated logs not available for {application_id!r}. "
|
|
f"The application may not be in FINISHED state, or "
|
|
f"yarn.log-aggregation-enable is false on the cluster."
|
|
)
|
|
if resp.status_code >= 400:
|
|
logger.error(f"YARN GET {url} -> {resp.status_code}: {resp.text[:500]}")
|
|
raise YarnError(f"YARN GET logs returned HTTP {resp.status_code}")
|
|
logger.info(f"YARN logs {application_id} -> {len(resp.text)} chars")
|
|
return resp.text
|
|
|
|
|
|
def kill_application(application_id: str, yarn_rm_url: str | None) -> None:
|
|
"""PUT state=KILLED to /ws/v1/cluster/apps/{appid}/state."""
|
|
url = f"{_base_url(yarn_rm_url)}/ws/v1/cluster/apps/{application_id}/state"
|
|
resp = _request("PUT", url, json_body={"state": "KILLED"})
|
|
if resp.status_code >= 400:
|
|
logger.error(f"YARN PUT {url} -> {resp.status_code}: {resp.text[:500]}")
|
|
raise YarnError(f"YARN kill returned HTTP {resp.status_code}: {resp.text}")
|
|
logger.info(f"YARN kill {application_id} -> ok")
|