Writing Order Data

Writing and Appending Order Data

The load reads JSON, Columnar and Binary Formats's orders.jsonl with an explicit schema, explodes each order's items into lines, and appends one year at a time, one snapshot per load as a nightly pipeline would. writeTo(...).append() is Spark 129 's DataFrameWriterV2 API, which commits through the catalog:

load_orders.py: two appends, then the cross-check against Chapters 2 and 4Python
"""Load Chapter 3's 100,000 sample orders into Iceberg in two appends: 2025, then 2026."""
from pyspark.sql import functions as F
from lake import spark
src = spark.read.json("/mnt/d/Books/Data Engineering/demos/ch03/data/orders.jsonl",
    schema="order_id BIGINT, customer_id INT, order_ts STRING, channel STRING, status STRING, "
           "coupon STRING, discount DECIMAL(9,2), total DECIMAL(10,2), "
           "items ARRAY<STRUCT<book_id: INT, qty: INT, unit_price: DECIMAL(9,2)>>")
src = src.withColumn("order_ts", F.to_timestamp("order_ts"))
orders = src.drop("items")
items = src.select("order_id", F.posexplode("items").alias("pos", "it")).select(
    "order_id", (F.col("pos") + 1).alias("line_no"), "it.book_id", "it.qty", "it.unit_price")
for year in (2025, 2026):                         # one commit (snapshot) per append
    o = orders.where(F.year("order_ts") == year)
    o.sortWithinPartitions("order_ts").writeTo("booknest.orders").append()
    items.join(o.select("order_id"), "order_id").writeTo("booknest.order_items").append()
spark.sql("""
  SELECT count(DISTINCT o.order_id) AS orders, count(*) AS lines,
         sum(i.qty * i.unit_price) FILTER (WHERE o.status <> 'cancelled') AS gross
  FROM booknest.orders o JOIN booknest.order_items i USING (order_id)""").show()
Output
+------+------+----------+
|orders| lines|     gross|
+------+------+----------+
|100000|137944|3303427.30|
+------+------+----------+

The lakehouse reproduces the baseline of XML and Its Toolchain and Analytical SQL and Data Warehouses exactly. Each append wrote one file per month (12, then 6): since Iceberg 1.2.0 129 the writer asks Spark to hash-distribute rows by partition first. Appends are not idempotent, so a production load uses MERGE INTO on order_id or partition overwrites (Idempotency and Safe Backfills).