A spark-submit Job

Writing and Running a spark-submit Job

A Spark 129 application is an ordinary Python file that builds its own session. Leave the master and resources out of the code, so the same file runs on a laptop and on a cluster with different spark-submit flags. Here is BookNest's nightly orders_daily job.

jobs/orders_daily.py: BookNest's nightly batch jobPython
"""BookNest nightly job: copies sold and revenue per day, genre and title.
Usage: spark-submit jobs/orders_daily.py OUTPUT_DIR"""
import sys
from pyspark.sql import SparkSession, functions as F
spark = (SparkSession.builder.appName("orders_daily")             # master: from spark-submit
         .config("spark.sql.session.timeZone", "UTC").getOrCreate())  # days in UTC
lines = spark.read.parquet("data/order_lines.parquet").where("status = 'delivered'")
books = spark.read.parquet("data/books.parquet").selectExpr("id AS book_id", "title", "genre")
daily = (lines.join(F.broadcast(books), "book_id")
         .groupBy(F.to_date("order_ts").alias("day"), "genre", "title")
         .agg(F.sum("qty").alias("copies"),
              F.sum(F.col("qty") * F.col("unit_price")).alias("revenue")))
daily.write.mode("overwrite").parquet(sys.argv[1])
out = spark.read.parquet(sys.argv[1])
print(f"orders_daily: {out.count():,} rows, {out.select('day').distinct().count()} days, "
      f"shuffle partitions {spark.conf.get('spark.sql.shuffle.partitions')}")
spark.stop()
Submitting the job in local mode
spark-submit --master "local[4]" --name orders_daily \
  jobs/orders_daily.py out/orders_daily
Output
orders_daily: 2,896 rows, 544 days, shuffle partitions 200

Arguments after the script name reach sys.argv. The job condensed 1.38 million order lines into 2,896 day, genre and title rows, still with the default 200 shuffle partitions. A failed job exits non-zero, which is what a scheduler such as Airflow 129 (Orchestration and Pipelines) needs.