feat(module-1): доработаны материалы модуля 1 и удалён legacy-ноутбук
- Зачем: - нужен воспроизводимый вход в курс и единое место хранения планов по модулям. - Что: - добавлены `plans/README.md`, living plan Модуля 1, `START_HERE.md` и канонический ноутбук `01_environment_and_smoke_test.ipynb`, а `cluster_smoke.py` расширен до reusable helper и CLI smoke test. - уточнены onboarding-материалы и окружение: добавлены Spark UI в `README.md`, ресурсы хоста и креды MinIO в `START_HERE.md`, использован `NB_GID` в `jupyter/Dockerfile`, в ноутбуке усилены самостоятельные задания и добавлены `cell id`, а пояснения в `cluster_smoke.py` переведены на русский для студентов. - обновлены `README.md`, `AGENTS.md` и archive howto, удалены устаревшие `spark-basic-test.ipynb` и `03_partitioning_and_schema_evolution.ipynb`. - Проверка: - `python3 -m py_compile src/spark/cluster_smoke.py src/spark/__init__.py`. - `docker compose build spark-master` и `docker compose build jupyter`. - `docker compose up -d`, `docker compose exec jupyter python3 -c "from spark.cluster_smoke import create_spark_session, run_cluster_smoke; spark=create_spark_session(app_name='module-01-validation'); print(run_cluster_smoke(spark)); spark.stop()"` и `docker compose exec jupyter jupyter nbconvert --to notebook --execute /opt/work/01_environment_and_smoke_test.ipynb --output-dir /tmp --output module1-validation-2.ipynb`.
This commit is contained in:
@@ -0,0 +1,3 @@
|
||||
from .cluster_smoke import create_spark_session, format_smoke_report, run_cluster_smoke
|
||||
|
||||
__all__ = ["create_spark_session", "format_smoke_report", "run_cluster_smoke"]
|
||||
@@ -1,12 +1,84 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from pyspark.sql import SparkSession
|
||||
|
||||
spark = (
|
||||
SparkSession.builder
|
||||
.appName("test-cluster")
|
||||
.master("spark://spark-master:7077") # важно: не local[*]
|
||||
.getOrCreate()
|
||||
)
|
||||
DEFAULT_MASTER_URL = "spark://spark-master:7077"
|
||||
DEFAULT_APP_NAME = "cluster-smoke"
|
||||
DEFAULT_ROW_COUNT = 1_000_000
|
||||
|
||||
spark.range(0, 1000000).groupBy().sum().show()
|
||||
|
||||
spark.stop()
|
||||
def create_spark_session(
|
||||
app_name: str = DEFAULT_APP_NAME,
|
||||
master_url: str = DEFAULT_MASTER_URL,
|
||||
) -> SparkSession:
|
||||
"""Создать SparkSession для проверок окружения в Модуле 1.
|
||||
|
||||
Функция специально остаётся минимальной и опирается на смонтированный
|
||||
в контейнер `spark-defaults.conf` для остальной Spark/Iceberg-конфигурации.
|
||||
"""
|
||||
return SparkSession.builder.appName(app_name).master(master_url).getOrCreate()
|
||||
|
||||
|
||||
def run_cluster_smoke(spark: SparkSession, row_count: int = DEFAULT_ROW_COUNT) -> dict:
|
||||
spark_context = spark.sparkContext
|
||||
master_url = spark_context.master
|
||||
|
||||
if master_url.startswith("local"):
|
||||
raise RuntimeError(
|
||||
"Smoke test must run against the Spark cluster, not local mode."
|
||||
)
|
||||
|
||||
# Берём больше партиций, чем даёт defaultParallelism, чтобы задача точно
|
||||
# распределилась по worker-ам, а не схлопнулась в слишком маленький job.
|
||||
partition_count = max(spark_context.defaultParallelism * 2, 8)
|
||||
expected_sum = row_count * (row_count - 1) // 2
|
||||
actual_sum = (
|
||||
spark.range(0, row_count, 1, numPartitions=partition_count)
|
||||
.selectExpr("SUM(id) AS total_sum")
|
||||
.collect()[0]["total_sum"]
|
||||
)
|
||||
|
||||
if actual_sum != expected_sum:
|
||||
raise RuntimeError(
|
||||
f"Unexpected smoke result: expected {expected_sum}, got {actual_sum}."
|
||||
)
|
||||
|
||||
return {
|
||||
"app_name": spark_context.appName,
|
||||
"master_url": master_url,
|
||||
"spark_version": spark.version,
|
||||
"default_parallelism": spark_context.defaultParallelism,
|
||||
"partition_count": partition_count,
|
||||
"row_count": row_count,
|
||||
"result_sum": actual_sum,
|
||||
"expected_sum": expected_sum,
|
||||
}
|
||||
|
||||
|
||||
def format_smoke_report(report: dict) -> str:
|
||||
return "\n".join(
|
||||
[
|
||||
"Spark cluster smoke test passed.",
|
||||
f"app_name={report['app_name']}",
|
||||
f"master_url={report['master_url']}",
|
||||
f"spark_version={report['spark_version']}",
|
||||
f"default_parallelism={report['default_parallelism']}",
|
||||
f"partition_count={report['partition_count']}",
|
||||
f"row_count={report['row_count']}",
|
||||
f"result_sum={report['result_sum']}",
|
||||
f"expected_sum={report['expected_sum']}",
|
||||
]
|
||||
)
|
||||
|
||||
|
||||
def main() -> None:
|
||||
spark = create_spark_session()
|
||||
try:
|
||||
report = run_cluster_smoke(spark)
|
||||
print(format_smoke_report(report))
|
||||
finally:
|
||||
spark.stop()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
|
||||
Reference in New Issue
Block a user