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 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).