import json import threading import time import urllib.request from pyspark.sql import SparkSession spark = SparkSession.builder.appName("sg-dynamic-allocation").getOrCreate() sc = spark.sparkContext sc.setLogLevel("ERROR") conf_keys = [ "spark.dynamicAllocation.enabled", "spark.dynamicAllocation.shuffleTracking.enabled", "spark.dynamicAllocation.minExecutors", "spark.dynamicAllocation.initialExecutors", "spark.dynamicAllocation.maxExecutors", "spark.executor.cores", ] for key in conf_keys: print(f"{key}={sc.getConf().get(key, '')}") app_id = sc.applicationId print(f"application_id={app_id}") def executor_ids(): url = f"{sc.uiWebUrl}/api/v1/applications/{app_id}/executors" try: with urllib.request.urlopen(url, timeout=5) as response: rows = json.load(response) except Exception as exc: print(f"executor_api_poll=retry reason={type(exc).__name__}") return [] return [row["id"] for row in rows if row["id"] != "driver" and row.get("isActive", True)] def slow_task(value): time.sleep(3) return value result = {} def run_job(): result["rows"] = sc.parallelize(range(48), 48).map(slow_task).count() worker = threading.Thread(target=run_job) worker.start() samples = [] for _ in range(18): ids = executor_ids() samples.append(len(ids)) print(f"active_executor_count={len(ids)} executor_ids={','.join(ids) if ids else 'none'}") if len(set(samples)) >= 2 and max(samples) >= 2: break time.sleep(1) worker.join() for _ in range(10): ids = executor_ids() samples.append(len(ids)) print(f"active_executor_count={len(ids)} executor_ids={','.join(ids) if ids else 'none'}") if samples[-1] < max(samples): break time.sleep(2) print("executor_count_samples=" + ",".join(str(value) for value in samples)) print(f"rows_processed={result['rows']}") spark.stop()