Spark SQL Deep Dive
Back in post 1, the nightly revenue report scanned a 2 TB orders table and took four hours. Here's the part nobody said in that meeting: the report selects three columns and filters on last week's dates. The query asks for a thimbleful; the engine drinks the ocean. Most of those four hours isn't computing revenue — it's reading bytes the answer never needed.
This post is about not reading the ocean. Spark SQL's speed on analytical queries comes less from clever computation than from aggressive avoidance: an optimizer (Catalyst) that rewrites your query to skip data, a file format (Parquet) that tells the reader exactly which chunks can be skipped, and a partitioning model that decides how the remaining work splits across machines. You'll learn the four stages every query passes through, what predicate pushdown and column pruning actually do and why they matter, how to read EXPLAIN output like a plan instead of a wall of text, and why partition keys and skew decide whether your cluster finishes in minutes or hours. The centerpiece is a real lab: a pure-Java program that writes a Parquet file, reads its footer statistics, and proves — from real bytes — that an engine can skip whole row groups without reading them.
Lab honesty, up front. Spark's SQL/DataFrame engine cannot run in this sandbox: its Netty transport dies with a "too large frame" error on every SQL operation (the same sandbox limitation documented in post 1). So the Catalyst, pushdown, EXPLAIN, and partitioning sections below describe standard documented behavior and are each marked [NOT RUN] — no invented query plans, no invented outputs anywhere. What did run is the Parquet lab: pure Java against the parquet-mr 1.13.1 libraries bundled with Spark 3.5.3 (no Spark engine, no network), writing a real Parquet file, reading its real footer metadata, and computing real skip decisions. That lab is the proof; everything else is described, not run.
Principle: the fastest I/O is the I/O you never do. Every optimization in this post is a variation on that sentence. Hold onto it and the rest reads as consequences.
Catalyst: your SQL becomes a plan in four stages
[NOT RUN] — Spark's SQL engine can't run in this lab (Netty transport failure); this section describes standard documented behavior.
Post 3 introduced the DataFrame API: named, typed columns instead of opaque JVM objects. The payoff for that design is Catalyst, Spark's query optimizer. Because your query is expressed in columns and operations the engine understands — not arbitrary Java lambdas — the engine can rewrite it before running it, the way a database rewrites SQL. Every query passes through four stages:
- Parsed plan (unresolved). The SQL text becomes a syntax tree. Column references are still just names — nothing is bound to a table yet. If your SQL has a syntax error, it dies here.
- Analyzed plan (resolved). The catalog binds every name to a real table and column and checks types. "Column not found" and type-mismatch errors surface here. After this stage the plan is resolved: every name means something.
- Optimized plan. Rule-based rewrites transform the plan into an equivalent cheaper one: pushing filters down toward the scan, pruning unneeded columns, folding constants, simplifying boolean expressions, propagating nulls. These are heuristics — good bets that usually win, not guarantees.
- Physical plan. The logical plan becomes executable operators: which join algorithm (broadcast the small side vs shuffle both sides), how to scan the files, how data is partitioned between stages. Spark can even adjust some of these choices at runtime (Adaptive Query Execution) — but the shape of the win is decided in stage 3.
One honest limitation, stated as a rule of thumb: Catalyst can only optimize what it can see. A predicate written in plain column expressions is transparent to the optimizer; the same logic hidden inside a UDF is an opaque black box — the optimizer can't push a filter it can't inspect. Prefer column expressions over UDFs where the logic allows it, not for style points, but because it keeps the query optimizable.
Predicate pushdown and column pruning: less I/O, on purpose
[NOT RUN] — Spark's SQL engine can't run in this lab (Netty transport failure); this section describes standard documented behavior.
Two of the optimizer's rules do most of the heavy lifting on analytical scans. They deserve precise definitions, because people wave the terms around loosely:
- Predicate pushdown moves filters as close to the data source as possible — ideally into the scan itself, so rows that can't match are never read, never decompressed, never decoded.
WHERE age > 30shouldn't filter a million rows in memory; it should prevent most of those rows from being read at all. - Column pruning makes the scan read only the columns the query references.
SELECT name, ageover a 50-column table shouldn't touch the other 48 columns' bytes.
Why do these two matter so much? Because of what analytical queries look like in practice: a few columns out of many, a slice of rows out of billions. On large scans, reading the bytes — from disk or over the network — is usually the dominant cost (rule of thumb, not a law: on tiny cached datasets the CPU dominates, and then these rules buy you little). Skipping bytes is skipping time, and these two rules are how the engine skips bytes before the expensive work starts.
The row-vs-columnar distinction is what makes pruning powerful, stated precisely: Parquet stores each column's values together, so pruning a column means those bytes are simply never read. Contrast a row-oriented store: Postgres has indexes and a capable optimizer, and a good plan avoids full scans — but the row layout means a selected row brings all its columns along. Pruning helps everywhere; it transforms the economics of columnar storage specifically. (Post 5 opens the format hood completely.)
But there's a catch the definitions hide: a filter can only skip data the scan can evaluate without reading it. Pushdown is the plan; it needs a mechanism — per-chunk statistics stored in the file itself, consulted before a single data page is touched. That mechanism is the key insight of this post, and it's where the lab comes in.
The key insight: the footer knows before the scan reads
A Parquet file is not a flat pile of rows. It is organized in layers, and the last thing written is the first thing a smart reader consults:
- Row groups — horizontal slices of rows (production default is on the order of 128 MB; our lab uses a tiny size on purpose). A row group is the unit of skipping: an engine reads or skips it whole.
- Column chunks — inside each row group, one chunk per column. This is the physical shape that makes column pruning possible: the
agebytes sit together, apart from thenamebytes. - Pages — inside each column chunk, the actual encoded values.
- The footer — written at the very end of the file: the schema, plus per-column-chunk statistics: the minimum value, the maximum value, the null count, and the encodings used.
Now the trick. The reader opens the file, reads the tiny footer first, and consults the statistics before touching any data page. Suppose the query filters age > 30 and the footer's statistics say row group 0 has min(age)=18, max(age)=29. No row in that group can possibly match — the whole row group is skipped: no read, no decompression, no decode. Row group 1 with max(age)=45 can't be ruled out, so it's read and evaluated row by row. That is predicate pushdown's file-level machinery: the optimizer pushes the filter into the scan, and the scan uses footer statistics to skip row groups.
Notice what this implies about data layout: pushdown works best when each row group's min/max range is narrow — which happens when rows with similar values sit near each other (sorted or clustered data). Randomly interleaved ages would give every row group a min of 18 and a max of 65, and no group could ever be skipped. File layout is a performance decision, not just a storage decision — a theme post 5 develops fully.
Lab: prove row-group skipping with real footer bytes
This ran — Temurin JDK 17, the parquet-mr 1.13.1 jars bundled with Spark 3.5.3, no Spark engine, no network. Two small programs. The first writes a Parquet file with three deliberately separated age bands, one per row group:
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.parquet.example.data.Group;
import org.apache.parquet.example.data.simple.SimpleGroupFactory;
import org.apache.parquet.hadoop.example.ExampleParquetWriter;
import org.apache.parquet.hadoop.example.GroupWriteSupport;
import org.apache.parquet.hadoop.ParquetWriter;
import org.apache.parquet.hadoop.metadata.CompressionCodecName;
import org.apache.parquet.io.api.Binary;
import org.apache.parquet.schema.LogicalTypeAnnotation;
import org.apache.parquet.schema.MessageType;
import org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName;
import org.apache.parquet.schema.Types;
public class WritePeople {
static final MessageType SCHEMA = Types.buildMessage()
.required(PrimitiveTypeName.INT32).named("age")
.required(PrimitiveTypeName.BINARY).as(LogicalTypeAnnotation.stringType()).named("name")
.optional(PrimitiveTypeName.BINARY).as(LogicalTypeAnnotation.stringType()).named("city")
.named("Person");
public static void main(String[] args) throws Exception {
long rowGroupSize = args.length > 0 ? Long.parseLong(args[0]) : 128 * 1024 * 1024L;
Configuration conf = new Configuration();
GroupWriteSupport.setSchema(SCHEMA, conf);
conf.setBoolean("parquet.enable.dictionary", false);
Path out = new Path("/tmp/parquetlab/people.parquet");
org.apache.hadoop.fs.FileSystem.get(conf).delete(out, false);
ParquetWriter<Group> writer = ExampleParquetWriter.builder(out)
.withConf(conf)
.withRowGroupSize(rowGroupSize)
.withDictionaryEncoding(false)
.withCompressionCodec(CompressionCodecName.SNAPPY)
.build();
SimpleGroupFactory gf = new SimpleGroupFactory(SCHEMA);
int id = 0;
// Band A: ages 18-29 (100 rows)
for (int i = 0; i < 100; i++) { write(writer, gf, 18 + (i % 12), "user-a-" + i, "springfield"); id++; }
// Band B: ages 30-45 (100 rows)
for (int i = 0; i < 100; i++) { write(writer, gf, 30 + (i % 16), "user-b-" + i, "shelbyville"); id++; }
// Band C: ages 46-65 (100 rows, city null every other row)
for (int i = 0; i < 100; i++) { write(writer, gf, 46 + (i % 20), "user-c-" + i, (i % 2 == 0) ? null : "ogdenville"); id++; }
writer.close();
System.out.println("wrote " + id + " rows, rowGroupSize=" + rowGroupSize);
}
static void write(ParquetWriter<Group> w, SimpleGroupFactory gf, int age, String name, String city) throws Exception {
Group g = gf.newGroup();
g.append("age", age);
g.append("name", Binary.fromString(String.format("%-40s", name)));
if (city != null) g.append("city", Binary.fromString(city));
w.write(g);
}
}
Three bands of 100 rows with disjoint age ranges — 18–29, 30–45, 46–65 — and an optional city column that's null every other row in band C, so the footer has null counts to show. The 4500-byte row-group size was found by trial: it's just under one band's buffered size, so each band flushes as exactly its own row group. (Production row groups are ~128 MB — this is a lab dial, not a recommendation.) Dictionary encoding is off so the encodings stay plain and readable; Snappy compression stays on, because real files are compressed.
The second program reads the footer and plays the role of the scan: it prints every row group's per-column statistics, then applies the filter age > 30 to those statistics exactly the way an engine would — skip the group only if max <= threshold — and finally computes the column-pruning math from the real compressed byte counts:
import java.util.*;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.parquet.hadoop.ParquetFileReader;
import org.apache.parquet.hadoop.metadata.*;
public class ReadFooter {
static String render(Object v) {
if (v instanceof org.apache.parquet.io.api.Binary) {
return "'" + ((org.apache.parquet.io.api.Binary) v).toStringUsingUTF8().trim() + "'";
}
return String.valueOf(v);
}
public static void main(String[] args) throws Exception {
Configuration conf = new Configuration();
Path in = new Path("/tmp/parquetlab/people.parquet");
ParquetMetadata meta;
try (ParquetFileReader reader = ParquetFileReader.open(conf, in)) {
meta = reader.getFooter();
}
FileMetaData fileMeta = meta.getFileMetaData();
List<BlockMetaData> blocks = meta.getBlocks();
System.out.println("file: " + in);
System.out.println("schema: " + fileMeta.getSchema());
System.out.println("row groups: " + blocks.size());
System.out.println();
for (int b = 0; b < blocks.size(); b++) {
BlockMetaData block = blocks.get(b);
System.out.println("== row group " + b + ": " + block.getRowCount() + " rows, "
+ block.getTotalByteSize() + " bytes total ==");
for (ColumnChunkMetaData col : block.getColumns()) {
String path = String.join(".", col.getPath().toArray());
org.apache.parquet.column.statistics.Statistics<?> stats = col.getStatistics();
StringBuilder sb = new StringBuilder();
sb.append(String.format(" col %-6s rows=%-4d nulls=%-4d", path, col.getValueCount(),
stats == null ? -1 : stats.getNumNulls()));
if (stats != null && stats.hasNonNullValue()) {
sb.append(" min=").append(render(stats.genericGetMin())).append(" max=").append(render(stats.genericGetMax()));
} else if (stats == null || stats.isEmpty()) {
sb.append(" (no stats)");
}
sb.append(" encodings=").append(col.getEncodings());
sb.append(" compressed=").append(col.getTotalSize()).append("B");
sb.append(" uncompressed=").append(col.getTotalUncompressedSize()).append("B");
System.out.println(sb);
}
}
// --- the pushdown decision: filter age > threshold (default 30) ---
int threshold = args.length > 0 ? Integer.parseInt(args[0]) : 30;
System.out.println();
System.out.println("== predicate pushdown simulation: WHERE age > " + threshold + " ==");
long skippedBytes = 0, readBytes = 0;
int skippedGroups = 0;
for (int b = 0; b < blocks.size(); b++) {
BlockMetaData block = blocks.get(b);
ColumnChunkMetaData ageCol = block.getColumns().stream()
.filter(c -> String.join(".", c.getPath().toArray()).equals("age")).findFirst().get();
org.apache.parquet.column.statistics.Statistics<?> s = ageCol.getStatistics();
long lo = ((Number) s.genericGetMin()).longValue();
long hi = ((Number) s.genericGetMax()).longValue();
long compressedBytes = block.getColumns().stream().mapToLong(ColumnChunkMetaData::getTotalSize).sum();
String verdict;
if (hi <= threshold) { verdict = "SKIP"; skippedBytes += compressedBytes; skippedGroups++; }
else { verdict = "READ"; readBytes += compressedBytes; }
System.out.println(String.format(" row group %d: age in [%d, %d] -> max(%d) %s %d -> %s (%d compressed bytes)",
b, lo, hi, hi, (hi <= threshold ? "<=" : ">"), threshold, verdict, compressedBytes));
}
System.out.println(" -> would SKIP " + skippedGroups + " of " + blocks.size() + " row groups: "
+ skippedBytes + " compressed bytes never read; READ " + readBytes + " compressed bytes");
// --- column pruning: only age + name needed ---
System.out.println();
System.out.println("== column pruning: SELECT age, name (city dropped) ==");
Set<String> wanted = new HashSet<>(Arrays.asList("age", "name"));
long allBytes = 0, wantedBytes = 0;
for (BlockMetaData block : blocks) {
for (ColumnChunkMetaData col : block.getColumns()) {
String p = String.join(".", col.getPath().toArray());
allBytes += col.getTotalSize();
if (wanted.contains(p)) wantedBytes += col.getTotalSize();
}
}
System.out.println(" all columns: " + allBytes + " bytes compressed; age+name only: " + wantedBytes
+ " bytes (" + (100 * wantedBytes / allBytes) + "% of the data)");
}
}
Compile against the jars bundled with Spark 3.5.3, write, then read the footer — every line below is verbatim output:
$ export JAVA_HOME=~/workspace/tooling/jdk17
$ CP=$(ls ~/workspace/tooling/spark-3.5.3-bin-hadoop3/jars/*.jar | tr '\n' ':')
$ $JAVA_HOME/bin/javac -cp "$CP" -d /tmp/parquetlab WritePeople.java ReadFooter.java
[...one deprecation note about a deprecated API — compiles and runs fine...]
$ $JAVA_HOME/bin/java -cp "$CP:/tmp/parquetlab" WritePeople 4500
wrote 300 rows, rowGroupSize=4500
$ $JAVA_HOME/bin/java -cp "$CP:/tmp/parquetlab" ReadFooter
file: /tmp/parquetlab/people.parquet
schema: message Person {
required int32 age;
required binary name (STRING);
optional binary city (STRING);
}
row groups: 3
== row group 0: 100 rows, 6385 bytes total ==
col age rows=100 nulls=0 min=18 max=29 encodings=[BIT_PACKED, PLAIN] compressed=98B uncompressed=426B
col name rows=100 nulls=0 min='user-a-0' max='user-a-99' encodings=[BIT_PACKED, PLAIN] compressed=556B uncompressed=4426B
col city rows=100 nulls=0 min='springfield' max='springfield' encodings=[RLE, BIT_PACKED, PLAIN] compressed=123B uncompressed=1533B
== row group 1: 100 rows, 6385 bytes total ==
col age rows=100 nulls=0 min=30 max=45 encodings=[BIT_PACKED, PLAIN] compressed=116B uncompressed=426B
col name rows=100 nulls=0 min='user-b-0' max='user-b-99' encodings=[BIT_PACKED, PLAIN] compressed=556B uncompressed=4426B
col city rows=100 nulls=0 min='shelbyville' max='shelbyville' encodings=[RLE, BIT_PACKED, PLAIN] compressed=123B uncompressed=1533B
== row group 2: 100 rows, 5595 bytes total ==
col age rows=100 nulls=0 min=46 max=65 encodings=[BIT_PACKED, PLAIN] compressed=125B uncompressed=426B
col name rows=100 nulls=0 min='user-c-0' max='user-c-99' encodings=[BIT_PACKED, PLAIN] compressed=556B uncompressed=4426B
col city rows=100 nulls=50 min='ogdenville' max='ogdenville' encodings=[RLE, BIT_PACKED, PLAIN] compressed=85B uncompressed=743B
== predicate pushdown simulation: WHERE age > 30 ==
row group 0: age in [18, 29] -> max(29) <= 30 -> SKIP (777 compressed bytes)
row group 1: age in [30, 45] -> max(45) > 30 -> READ (795 compressed bytes)
row group 2: age in [46, 65] -> max(65) > 30 -> READ (766 compressed bytes)
-> would SKIP 1 of 3 row groups: 777 compressed bytes never read; READ 1561 compressed bytes
== column pruning: SELECT age, name (city dropped) ==
all columns: 2338 bytes compressed; age+name only: 2007 bytes (85% of the data)
Read that output the way an engine would. The footer gives each row group real min/max statistics for age: 18–29, 30–45, 46–65. Against the filter age > 30, the decision is mechanical: row group 0's max(29) <= 30, so no row in it can match — SKIP, 777 compressed bytes never read. Row group 1's range straddles the threshold, so it must be READ and evaluated row by row (rows aged 31+ match). Row group 2's min(46) > 30 means every row matches, but the engine still has to READ it to return those rows — statistics can only ever rule out a row group, never certify its contents. Notice the footer reports each row group's total as uncompressed bytes (6385), while the skip math uses compressed bytes (777/795/766) — the compressed numbers are what would actually travel off disk.
Two more details worth pausing on. First, city in row group 2 shows nulls=50 — the footer really does track null counts per column chunk (a newer Parquet feature; parquet-mr 1.13 writes them), and engines use them the same way: a filter like city IS NOT NULL can skip a row group where every value is null. Second, encodings=[BIT_PACKED, PLAIN]: PLAIN is the value encoding (dictionary encoding was disabled for readability); BIT_PACKED covers the repetition/definition level streams. Post 5 goes inside encodings properly.
The column-pruning line is the quiet one: SELECT age, name touches 2007 of 2338 compressed bytes — dropping city saves 331 bytes, 15%. On this 3-column toy that's modest; the principle is what scales. On a 50-column table where the query needs 3 columns, pruning skips the overwhelming majority of the bytes — and Catalyst applies it automatically from the query text, which is why writing SELECT * in analytics code is a habit worth breaking.
Reading EXPLAIN output: the plan, section by section
[NOT RUN] — Spark's SQL engine can't run in this lab (Netty transport failure); this section describes standard documented behavior.
EXPLAIN FORMATTED prints the query plan in four sections that mirror the Catalyst stages — and each section answers a different debugging question. I'm not pasting a plan here, because no engine ran to produce one; instead, here is the vocabulary for reading any real plan, section by section:
- Parsed Logical Plan — the query as written: a
Projectover aFilterover an unresolved relation. If something looks wrong here, the SQL text itself is wrong. - Analyzed Logical Plan — names bound to the catalog (
people.age, not justage), types resolved. Surprise casts and renamed columns show up here. - Optimized Logical Plan — the rules have run. Two things to look for: your
Filternow sits directly on the scan (that's predicate pushdown, visible in the plan), and theProjectlists only the columns the query uses (that's column pruning). If the filter is still floating above a projection, something blocked the pushdown. - Physical Plan — the executable version. The scan node names the format and carries the two fields worth reading:
PushedFilters, the predicates the file reader applies while scanning (this is whereGreaterThan(age,30)belongs — and where the footer statistics from the lab get consulted), andReadSchema, the columns actually read (justnameandageif pruning worked).
The node vocabulary beyond the scan: Filter is a filter that did not get pushed — ask why. Project is column selection. Exchange is a shuffle: data moving across the network between stages so equal keys land together — a rule of thumb is that every Exchange is a place where partitioning and skew can bite (next section). Sort and HashAggregate are what they sound like.
The diagnostic habit this builds: find the scan, read its PushedFilters. If your WHERE clause isn't there, the predicate was opaque to the optimizer — the usual suspect is a UDF or a non-deterministic expression (rule of thumb, not a law, but it's the first place to look). A plan where the filter sits on the scan and the read schema is minimal is a plan doing the least I/O it can; everything else is tuning.
Partitioning: how Spark splits the work
[NOT RUN] — Spark's SQL engine can't run in this lab (Netty transport failure); this section describes standard documented behavior.
Pushdown decides which bytes get read; partitioning decides how the work splits across machines. Partitions belong early in your mental model, because every performance story in Spark ends here:
- A partition is the unit of parallelism. As a rule of thumb, one task processes one partition — so the partition count is roughly your parallelism ceiling. DataFrames read from files get their partitions from file splits and row groups; operations like joins and aggregations create new partitioning via shuffle.
- The partition key decides the layout. A shuffle groups rows so equal keys land on the same partition — that's what makes a join or a
GROUP BYcorrect. The key you group or join by isn't just logic; it's a physical data-layout decision. Choose it deliberately. - Too few partitions and cores sit idle; too many and the scheduler drowns in tiny tasks. Spark's own tuning docs suggest on the order of 2–3 tasks per CPU core as a starting point — a rule of thumb, not a law. The right number is measured on your data, not copied from a blog.
Then there's the failure mode that no partition count fixes: skew. If one key holds a wildly disproportionate share of rows — one celebrity user with ten million events, one null key absorbing every malformed record — its partition dwarfs the others. The symptom is unmistakable: 199 tasks finish in a minute, one task runs for an hour, and the job's runtime is that one task. The standard remedy is salting: add a random prefix to the hot key, do a partial aggregate per salted key, then strip the prefix and aggregate again — spreading the hot key's rows across many partitions for the heavy step.
One more pruning flavor before the cheat sheet, because it completes the picture: partition pruning is the directory-level cousin of row-group skipping. With data laid out as events/year=2026/month=10/day=06, a filter on year and month lets Spark skip whole directories without even listing them — the directory names are the statistics. Row groups skip within a file; partitions skip across files. Post 5 covers partitioning schemes — and when this layout helps vs hurts — in depth.
Cheat sheet
- The fastest I/O is the I/O you never do. Every rule below is that sentence in disguise.
- Catalyst's four stages: parsed (syntax tree, names unresolved) → analyzed (names bound, types checked) → optimized (rules rewrite) → physical (join choice, scan, execution).
- Predicate pushdown moves filters into the scan so non-matching data is never read; column pruning makes the scan read only referenced columns. On large scans, skipped bytes are skipped time (rule of thumb, not a law).
- The file-level machinery: Parquet footers carry per-column-chunk min/max, null counts, and encodings. The reader consults them before touching data pages and skips whole row groups whose
[min, max]can't satisfy the filter. - Statistics can only rule out a row group, never certify it — a straddling range must still be read.
- Pushdown works best on clustered data: narrow min/max ranges per row group. Randomly interleaved values make every row group unskippable.
EXPLAIN FORMATTEDshows all four stages. Diagnostic habit: find the scan, readPushedFiltersandReadSchema. A filter that didn't push is usually opaque to the optimizer — suspect a UDF first.- A partition is the unit of parallelism (roughly one task per partition). Too few → idle cores; too many → scheduler overhead. Spark's docs suggest ~2–3 tasks per core as a starting point — measure, don't memorize.
- Skew: one hot key → one giant partition → one straggler task sets the job's runtime. Symptom: most tasks finish fast, one doesn't. Remedy: salting.
- Partition pruning skips whole directories (
year=2026/month=10/…) the way footer stats skip row groups within a file. - Postgres has indexes and a real optimizer — "reads every row" is never the right mental model. Row-vs-columnar is about layout: a selected row brings all its columns along; a selected column brings only its bytes.
What's next
You now know why skipping works — the optimizer rules, the footer statistics, the partition layout. But the format underneath is still a black box: what actually lives inside a column chunk, how encodings like dictionary and run-length encoding shrink the bytes, why Snappy is the usual default, and how partitioning schemes and Delta Lake turn a pile of Parquet files into something you can trust with concurrent writes and schema changes.
The next post, File Formats & the Lakehouse, opens the hood: Parquet pages and encodings from real bytes, partitioning schemes and when they hurt, and Delta Lake's transaction log — the machinery that makes the lakehouse more than a directory of files.
Field check before you move on: (1) Re-run the lab with your own filter — pass a threshold as an argument (e.g. ReadFooter 25 for age > 25). Predict which row groups skip before you run it, then verify. (2) Add a fourth band (ages 66–80, 100 rows) to WritePeople, re-run, and predict the skip set for age > 30 from the footer before looking. (3) Edit the wanted set in ReadFooter to city only, re-run, and compute how many compressed bytes that scan would touch versus SELECT * — then explain why the saving is larger than the age, name case. (4) Paper exercise: sketch the four Catalyst stages for SELECT name FROM people WHERE age > 30 — which stage pushes the filter into the scan, and which footer statistics make that pushdown effective?
Continue: Java Learning Roadmap 2026
Comments
Post a Comment