# coding=utf-8 """ @Time :2026/7/8 @Author :tao.chen Tools for inspecting YARN applications NOT submitted through this service. These bypass the local JobStore and require the caller to supply both the YARN application_id and the name of a saved Connection. """ import json from common.logging import logger from spark_executor.core.yarn_client import ( YarnClientConfig, get_application_logs, get_application_status, list_applications as list_applications_yarn, # alias to avoid collision ) from spark_executor.models import ApplicationSummary, JobResult, JobStatus from spark_executor.tools.connections import store as conn_store def get_external_job_logs(application_id: str, connection_name: str, tail_chars: int = 5000) -> str: """Fetch aggregated container logs for a YARN application by ID.""" logger.debug(f"get_external_job_logs enter application_id={application_id} connection_name={connection_name} tail_chars={tail_chars}") conn = conn_store.get(connection_name) if conn is None: raise KeyError(f"Connection not found: {connection_name}") config = YarnClientConfig.from_connection(conn) full = get_application_logs(application_id, config) tailed = full[-tail_chars:] if len(full) > tail_chars else full logger.info( f"get_external_job_logs ok application_id={application_id} connection_name={connection_name} " f"full_chars={len(full)} returned_chars={len(tailed)}" ) return tailed def get_external_job_status(application_id: str, connection_name: str) -> JobStatus: """Query YARN for an external application's current status.""" logger.debug(f"get_external_job_status enter application_id={application_id} connection_name={connection_name}") conn = conn_store.get(connection_name) if conn is None: raise KeyError(f"Connection not found: {connection_name}") config = YarnClientConfig.from_connection(conn) state, raw = get_application_status(application_id, config) logger.info(f"get_external_job_status ok application_id={application_id} connection_name={connection_name} state={state}") return JobStatus(application_id=application_id, state=state, raw=raw) def get_external_job_result(application_id: str, connection_name: str) -> JobResult: """Query YARN for an external application's terminal result view.""" logger.debug(f"get_external_job_result enter application_id={application_id} connection_name={connection_name}") conn = conn_store.get(connection_name) if conn is None: raise KeyError(f"Connection not found: {connection_name}") config = YarnClientConfig.from_connection(conn) state, raw = get_application_status(application_id, config) app = json.loads(raw).get("app", {}) result = JobResult( application_id=application_id, state=state, final_status=app.get("finalStatus"), diagnostics=app.get("diagnostics"), tracking_url=app.get("trackingUrl"), started_time=app.get("startedTime"), finished_time=app.get("finishedTime"), ) logger.info(f"get_external_job_result ok application_id={application_id} connection_name={connection_name} state={state}") return result def list_applications( connection_name: str, state: str | None = None, queue: str | None = None, limit: int = 100, ) -> list[ApplicationSummary]: """List YARN applications on the named cluster, optionally filtered. Bypasses the local JobStore (this is for apps not submitted through this service). The YARN ResourceManager REST endpoint /ws/v1/cluster/apps is queried through the connection's auth/SSL config. Defaults: limit=100 (YARN has no offset-based pagination, so large clusters should use state/queue filters to scope the result). """ logger.debug( f"list_applications enter connection_name={connection_name} " f"state={state} queue={queue} limit={limit}" ) conn = conn_store.get(connection_name) if conn is None: raise KeyError(f"Connection not found: {connection_name}") config = YarnClientConfig.from_connection(conn) raw_apps = list_applications_yarn( config, state=state, queue=queue, limit=limit ) summaries = [ ApplicationSummary( application_id=app.get("id", ""), name=app.get("name", ""), user=app.get("user", ""), queue=app.get("queue", ""), state=app.get("state", ""), final_status=app.get("finalStatus"), application_type=app.get("applicationType"), application_tags=app.get("applicationTags", ""), started_time=app.get("startedTime", 0), finished_time=app.get("finishedTime", 0), tracking_url=app.get("trackingUrl"), progress=app.get("progress"), ) for app in raw_apps ] logger.info( f"list_applications ok connection_name={connection_name} count={len(summaries)}" ) return summaries