Loading Order History

Loading BookNest's Order History

The pieces combine into the job that turns BookNest's raw order history (one million generated sample orders as JSON Lines) into typed Parquet 129 tables. It reads with the explicit schema, keeps any unparseable line in a _bad column instead of losing it, checks the data, and only then writes. It ran as spark-submit jobs/load_history.py data/raw out/history with the workstation's spark-defaults.conf (Application Configuration).

jobs/load_history.py: loading and checking BookNest's order historyPython
"""Load BookNest's raw order history (JSON Lines) into typed, checked Parquet tables.
Usage: spark-submit jobs/load_history.py RAW_DIR OUT_DIR"""
import sys
from pyspark.sql import SparkSession, functions as F
ORDER_SCHEMA = """order_id BIGINT, customer_id INT, order_ts TIMESTAMP, channel STRING,
  status STRING, currency STRING,
  items ARRAY<STRUCT<book_id: INT, qty: INT, unit_price: DECIMAL(6,2)>>,
  coupon STRING, discount DECIMAL(8,2), total DECIMAL(10,2), _bad STRING"""
raw, out = sys.argv[1], sys.argv[2]
spark = SparkSession.builder.appName("load_history").getOrCreate()
orders = (spark.read.schema(ORDER_SCHEMA).option("columnNameOfCorruptRecord", "_bad")
          .json(f"{raw}/orders.jsonl"))
gross = F.aggregate("items", F.lit(0).cast("decimal(10,2)"),
                    lambda acc, i: (acc + i.qty * i.unit_price).cast("decimal(10,2)"))
c = orders.agg(F.count("*").alias("orders"), F.count("_bad").alias("malformed"),
               F.count_if(gross - F.col("discount") != F.col("total")).alias("bad_totals"),
               F.sum("total").alias("revenue")).first()
print(f"{c.orders:,} orders, {c.malformed} malformed, {c.bad_totals} bad totals, "
      f"revenue {c.revenue:,}")
if c.malformed or c.bad_totals:
    sys.exit("load_history: checks failed, nothing written")
orders = orders.drop("_bad")
orders.write.mode("overwrite").parquet(f"{out}/orders")
lines = orders.select("order_id", "customer_id", "order_ts", "channel", "status",
                      F.inline("items"))
lines.write.mode("overwrite").parquet(f"{out}/order_lines")
print(f"wrote orders and {spark.read.parquet(f'{out}/order_lines').count():,} order lines")
spark.stop()
Output
1,000,000 orders, 0 malformed, 0 bad totals, revenue 34,794,290.24
wrote orders and 1,380,694 order lines

All checks share one aggregation, so one pass over the file serves them all. bad_totals recomputes each total from its lines in exact decimal arithmetic; zero mismatches and the 34,794,290.24 revenue match the generator's figures. A failed check exits non-zero before anything is written, so a scheduler (Orchestration and Pipelines) sees it and the previous tables survive. F.inline turns the items structs into order-line rows. Spark SQL and Joins to Partitioning and Caching read the same tables from data/ as orders, lines, customers and books.