Guides

The posts on this site in the order worth reading them, grouped by what you are trying to learn: how Spark executes a query, how the table formats store a table, and which tool to reach for.

The category and archive pages list everything by topic and by date. This page is different: it puts the posts in a reading order, and says what each one gives you, so you can start somewhere sensible rather than at whatever was published most recently.

Each path is built so that every post assumes only what the ones above it already explained.

How does Spark actually execute a query?

Ten posts that build on each other, from the runtime up to the optimizer and back down to what fails in production. If you read one, read the first.

  1. Apache Spark architecture is the foundation: who schedules what, how a job becomes stages and tasks, where a shuffle is written, and which component to suspect when a job misbehaves. Everything else here assumes it.
  2. Inside Spark’s Catalyst optimizer takes one query apart rule by rule: the four trees, what a rule is in the source, and how to switch one off and watch the plan change.
  3. Adaptive query execution covers what the engine revises once it has measured the data, and, just as usefully, the parts of the plan it never revisits.
  4. Apache Spark joins in depth is the decision that dominates most jobs: five strategies, why there are five, and how to make each one appear on demand.
  5. Spark memory management accounts for every region of an executor’s memory and the arithmetic that sizes it, including the one hard floor you cannot tune below.
  6. Spark shuffle internals goes under the stage boundary the first post described: Spark has three shuffle writers and picks one per shuffle without ever naming the choice. The decision tree here is read from the source and confirmed by running it.
  7. Data skew in Apache Spark is the most common way all of the above goes wrong in production, with six fixes and the measurements that tell you which one you need.
  8. Spark performance tuning is the capstone: seventeen techniques, each with the code, the evidence the optimization actually fired, and the measurement that says whether it helped. Most tuning advice is a list of settings with no way to tell.
  9. Accelerating Spark with DataFusion Comet swaps those operators for native ones without changing the query, which is the clearest way to see which parts of a plan the JVM was costing you.
  10. Every Apache Spark release puts the whole thing on a timeline: what each version changed, which defaults moved underneath you, and what a given upgrade is actually worth.

What changes once the query never ends?

Streaming is the same engine with a different contract: the query outlives any one run, so correctness depends on what it wrote down before it stopped.

  • Structured Streaming internals is the place to start: a streaming query is a loop that writes what it is about to do, does it, then writes that it finished. This is what that loop puts on disk, how the state store is laid out underneath it, and which parts of a query a restart can never change.
  • Streaming into Iceberg is the same loop pointed at a table format, where every micro-batch commits a snapshot and the trigger you pick decides how fast your metadata grows. It appears again as step 10 of the Iceberg path below, which is the other way into it.

For changing the log level on a query you cannot restart, see Logging a Spark stream in the reference list above.

Which Spark reference should I keep open?

These are lookup material rather than a reading order. Open them beside your editor.

How do the table formats store a table?

Start by choosing, then go deep on the one you chose.

  1. Open table formats in practice is the comparison: what actually differs between Iceberg, Hudi and Delta Lake, rather than what the marketing pages differ on.
  2. Apache Hudi, if you chose Hudi: it has its own reading path below, seven posts from the first upsert to concurrent writers.
  3. Apache Iceberg architecture for what a commit actually swaps, and why the catalog is part of the table.
  4. Apache XTable incremental sync for keeping one table readable as more than one format without re-converting everything each time.

The matching reference for step 4 is the XTable cheat sheet.

How does a Hudi table work, end to end?

Seven posts, from what an upsert does on disk to what goes wrong when two writers commit at once. Each one assumes only the ones above it.

Keep the Hudi cheat sheet open while you read: table types, the timeline, writer properties, indexes, table services and procedures, with the config key and default for each.

  1. Apache Hudi: an introduction is the place to start.
  2. Apache Hudi architecture for what happens when you upsert a row, down to bytes on object storage.
  3. Key generators decide what a record’s identity is and where it lands — the one Hudi decision you cannot change after the first write.
  4. The Hudi index is the single biggest lever on upsert cost — it decides which file group a record belongs to.
  5. Hudi Merge-on-Read for where the write you skipped goes, and which read pays for it.
  6. Incremental processing, the capability Hudi is named for: what an incremental read actually returns, how CDC mode differs, and the timestamp change in 1.x that makes the obvious checkpoint loop re-read every batch.
  7. Concurrency control, where two of three configurations lost a whole commit without raising anything, including the one most people deploy.

For running Hudi from a thin client or making a stalled write explain itself, see the Spark Connect and logging posts under How do I run a table format on Spark?

How does an Iceberg table work, end to end?

The longest path on the site, and the one that rewards being read in order: eighteen posts from why the format exists to what you have to schedule in production. Each one assumes only the ones above it.

Keep the Iceberg cheat sheet open while you read: the metadata tables, row-level operation modes, branches and tags, the CALL procedures and the maintenance jobs, with the property and default for each.

The shape of a table

  1. What Hive tables could not do is the motivation: a Hive table is a directory, and every limitation follows from that. Iceberg replaces the directory with a metadata tree that names every file.
  2. The three tiers — catalog, metadata, data — because knowing which layer a problem lives in is most of operating one.
  3. Iceberg catalogs answer one question, which metadata file is current, and how they answer it decides whether concurrent writers are safe.

How data moves through it

  1. Life of a write: three files and one pointer swap, and what ACID actually means here — measured with two writers hammering the same table.
  2. Life of a read: four prunes before a byte of data is touched, the last of which only works if your data is sorted.
  3. Format versions for what v1 to v4 change — chiefly whether a row can be deleted without rewriting a file.

Changing a table without rewriting it

  1. Schema evolution: why a rename costs nothing, and why field ids are the whole mechanism.
  2. Hidden partitioning and partition evolution: changing the layout without rewriting history.

Getting data in

  1. Appends, overwrites and row-level SQL, and the different mark each leaves in the snapshot log.
  2. Streaming ingestion, where the trigger you choose decides how fast your metadata grows.
  3. Copy-on-write or merge-on-read: the same DELETE producing two completely different file layouts.
  4. CDC into Iceberg, which works on day one and degrades quietly: three merges, three delete files, 960 dead records, and a row count that never moved.

Operating one

  1. Time travel and rollback: undoing a bad write in one statement.
  2. Migrating Hive tables — snapshot first, migrate second, and why that order matters.
  3. Sorting and clustering: what the per-file min/max statistics can and cannot skip, measured — the same query going from 48 files scanned to 1, and the z-order that made it worse.
  4. Running Iceberg in production: the settings fixed at create time, and the jobs nobody schedules.
  5. Iceberg views: a view definition that is versioned the way a table is — and a cross-engine promise that did not hold when it was tested between Spark and Trino.
  6. Branching, tagging and write-audit-publish: validating a write before anyone can read it, with the copy removed — plus the read redirect that catches people, and why there is no merge.

Two posts sit alongside this path rather than in it: table maintenance for Iceberg and Hudi, on why compaction reclaims nothing until a separate expiry job runs, and one table, three formats, on reading the same Parquet files as Iceberg, Delta and Hudi.

How do I run a table format on Spark?

The posts where the two halves of the site meet.

How do I run any of this locally?

Everything measured on this site was measured on a laptop. These are the environments that made that possible, in increasing order of how much they bring up.

  • A local Iceberg playground is the smallest thing that works: Spark, MinIO and a REST catalog in one compose file, when Iceberg is all you need.
  • Run every demo on this blog locally is the full stack — MinIO, a Hive Metastore, Spark, Kafka and Trino — plus the demo scripts that reproduce the numbers in these posts on your own machine.

Which tool do I reach for?

Small, self-contained utilities and command-line notes.

Everything, from basic to advanced