Pay Once, Not Per Task: Four Repeated Costs Comet 1.1.0 Stopped Paying

Apache DataFusion Comet 1.1.0 is out, with 14 of my PRs in it. Four of them fix the same kind of waste: work that gives the same answer every time but was being redone for every task, batch or shuffle block. Here is what each one changed, what the benchmarks really show, and the correctness fix I'd point to first.

· 10 min read
apache spark datafusion comet rust scala data engineering performance open source

Apache DataFusion Comet 1.1.0 is out: 398 commits from 40 contributors in the seven weeks since 1.0.0. The project’s release post walks through the big features, like native Iceberg writes, Celeborn remote shuffle and S3 credential handling, and the 1.1.0 changelog has everything else.

My last two posts were about reading less data: page indexes that EMR and Glue ignore, and a native Delta scan that respects them. This one is about a different kind of waste, one that has nothing to do with how much data you have. Four of my PRs in this release fix the same pattern in four places: work that comes out identical every time, yet was being redone for every task, every batch, or every shuffle block.

The Cost a Benchmark Run Hides

A small expense that recurs on every task is easy to overlook. A test run might pay it a few hundred times; a big production job can pay it thousands of times. At Spark’s default of 200 shuffle partitions and 8,192 rows per batch, a few hundred microseconds of setup per task barely registers. Move to 10,000 partitions, a table with 1,000 columns, or tiny batches left over after a selective filter, and that same overhead stacks up fast. No single query profile flags it. You just see executors that stay busy doing things nobody asked for.

For each PR below I asked three questions: what is being rebuilt, how often does that happen, and does the result ever actually change? (The regex, shuffle and zstd benchmarks all ran on an M-series Mac.)

Parse the Plan Once Per Stage (#5615)

Some quick background. Spark splits a job into stages, and each stage into many tasks, one per slice of the data. Comet sends every task its query plan as a compact binary message (protobuf). All the tasks in a stage receive exactly the same plan, apart from the list of files each one should read.

Even so, every task decoded the whole plan from scratch. It also rebuilt a lookup key for the scan by turning the schema and filters into text, which on wide tables turned out to be the slowest step of all, and it re-read the scan settings that every task shares (schema, filters, options). Only then did it add its own file list.

#5615 lets each executor hold on to the decoded plan and hand it to every later task in the stage. The driver computes a fingerprint of the plan once and ships it along, so executors can recognize a plan they’ve already seen without hashing it themselves. The scan’s lookup key is now worked out once on the driver too. And when the scan runs inside Comet’s native shuffle writer, all the map tasks of that shuffle share one decoded copy of the scan settings, which is discarded when the shuffle is cleaned up.

What I was most careful about is the piece that must not be shared. Each task’s file list really is different, so it’s still added task by task and never stored. A test checks that tasks get the same settings object while their file lists stay separate. A cache that leaked one task’s files into another would return wrong results, and no benchmark would ever notice.

Planning time per task (decoding, building the key and re-encoding, measured over 5,000 runs after warmup):

Scan planBeforeAfter
100 columns274 to 380 µs44 to 73 µs
1,000 columns2.0 to 2.5 ms0.55 to 0.93 ms

Computing the fingerprint on the driver shaved another 21% off the cached path for the 1,000-column plan (single thread, JDK 17). To put that in perspective, and this is arithmetic rather than a measurement: at those rates, a 1,000-column stage with 10,000 tasks would have burned 20 to 25 seconds of executor CPU on planning alone. Now it’s 5.5 to 9.3 seconds.

Compile the Regex Once Per Expression (#5612)

Before a regular expression can run, it has to be compiled into a matcher, and that step isn’t free. Comet’s rlike already did it once, when the query was planned. Three other functions, regexp_extract, regexp_extract_all and split, did it again for every batch of rows. They couldn’t simply borrow rlike’s approach because of how they’re wired up: Comet creates them by name, and they only see the pattern when they’re called.

#5612 gives each of these expressions room to remember one thing: its compiled pattern. The first batch compiles it, later batches reuse it, and if the pattern ever changes (split allows that) the new one gets compiled instead. Handing out a compiled regex is cheap, little more than bumping a reference count. Error messages are unchanged, and a bad pattern still fails at the same point as before.

In a benchmark, regexp_extract on standard 8,192-row batches dropped from 862 µs to 705 µs per batch, roughly 18% quicker. On 512-row batches it ran 2.1x faster. That’s the earlier point in miniature: the smaller the batch, the bigger the share of its time that went to compiling. The split numbers didn’t budge, because its test case used a plain text delimiter, which skips regex entirely. One open question I noted in the PR: the remembered pattern sits behind a lock, and the benchmarks ran on a single thread, so I haven’t measured what happens when many threads reach for it at once. Each lookup is a lock, a string comparison and a pointer copy, so I expect the effect to be negligible, but that’s an expectation, not a result.

Reuse Shuffle Scratch Per Partition (#5568)

A shuffle moves rows between stages, and each task writes its output split into many partitions. For every partition, the write path made two avoidable allocations. First, its output buffer started empty and grew step by step toward its 1 MB limit, so a 2,000-partition write went through that growth 2,000 times. Second, on every pass it copied the partition’s list of row positions into a new array twice the size, 16 bytes per row.

#5568 gives each task one buffer that it empties and refills partition after partition, both when spilling to disk and when finishing up. Row positions are now converted a small chunk at a time into a reusable scratch array rather than all at once. Tests compare the bytes written with a recycled buffer against a fresh one and confirm they match exactly.

With one task shuffling 4 million rows into 2,000 partitions, write time fell from 0.015 s to 0.011 s, about 27% less, with everything else unchanged. Those are small absolute numbers, and I’d rather admit that than dress them up. What the change really cuts is traffic to the memory allocator, which I expect to matter more when lots of tasks share it at once. A single-task benchmark can’t show that, and I haven’t measured it yet. Peak memory is the same as before; the difference is that a task now keeps exactly one write buffer no matter how many partitions it writes.

Reuse zstd Contexts Across Blocks (#5565)

Comet compresses shuffle data in blocks with zstd, and zstd needs a working area, called a context, to do its job. Every block created a fresh context and threw it away afterward. Setting one up is pure overhead, and it adds up fastest in shuffles with many partitions and small blocks.

#5565 keeps one context per task and reuses it for every block, resetting it and reapplying the compression level each time so a failure on one block can’t spill over into the next. Two details mattered. A reset context keeps whatever memory it grew to, so one unusually large block could leave it bloated for the rest of the task; an 8 MiB cap prevents that, and a test will fail if a future zstd upgrade pushes the level 8 context past that cap. And the remote shuffle path deliberately still creates a context per block, because its memory accounting reserves and releases that space around each one.

This is the most modest of the four wins, and I said as much in the PR. At Spark’s default 200 partitions the gain is lost in the noise, and in a one-task, 4-million-row test it was flat at 2,000 partitions. At 10,000 partitions and zstd level 3, over five runs, total time fell 5.6%, and the before and after runs formed ranges that didn’t overlap. Still, the gap between them is about the size of normal run-to-run variation, so compression time is the more trustworthy figure: 6.6% lower on the final run. A Linux run (arm64 in Docker, four dedicated CPUs) leaned the same way, about 3% overall and 6% on compression. An earlier version also reused the decompression context, but it didn’t speed anything up and held on to extra memory, so I took it out.

Fast Is Only Half the Job

Speed is worthless if the answers are wrong, and the PR from this release I’d point to first is #6041, a correctness fix. Comet’s native decimal SUM checked the running total against the column’s precision on every row and gave up the moment it overflowed. In the cases this fix covers, Spark lets the running total run past the precision and only checks when a value leaves its buffer, which is normally the final result. So summing the DECIMAL(38,38) values 0.6, 0.6, -0.6 gave 0.6 in Spark, since the middle step overflows but the final total fits, while Comet returned NULL, or an overflow error in ANSI mode (#6002). The fix makes Comet check at the same point Spark does, and at precision 38 it hands the query back to Spark whenever the planner can see the two would still disagree, such as a grouped sum under object hash aggregation or an aggregate Spark runs without code generation. The one case the planner can’t see, Spark abandoning code generation at runtime, is documented in the compatibility guide. As a bonus, dropping the per-row check sped up the ungrouped path: in a microbenchmark, ANSI mode went from 102 µs to 77 µs per 81,920 rows (legacy mode, 67 µs to 65 µs). The release also includes my fixes for matching how Spark resolves Parquet field IDs and duplicate field names (#5654, and #6116 through its backport #6266), and for turning away structs with duplicate field names before they can trip up Java Arrow (#5866). An accelerator that returns NULL where Spark returns a number isn’t faster. It’s just wrong sooner.

All 14 in 1.1.0

For the record, here’s everything of mine that shipped in 1.1.0, in the changelog’s own wording.

PRWhat it does
#5615Cache parsed plan data across a stage’s tasks
#5612Compile user regex patterns once per planned expression
#5568Reuse per-partition scratch in the shuffle write path
#5565Reuse zstd compression contexts across shuffle blocks
#6041Let decimal SUM recover from an intermediate overflow like Spark
#5654Match Spark’s duplicate field and field id semantics in Parquet field lookup
#6266Reject a file without field ids at any depth whether or not id matching is on (backport of #6116)
#5866Decline structs with duplicate field names before they reach Java Arrow
#6042Gate the regr_r2 degenerate-case swap on the Spark patch release
#5653Expand object store option references, uniquify constant metadata names, drop dead parquet JNI
#5874Run length on binary input natively
#5873Drop redundant width_bucket shim registrations, guard serde uniqueness
#6040Add a test guarding page skipping in the native Iceberg scan
#5762Label pull requests by changed paths and title prefix

How to Find These in Your Own Code

All four fixes share a shape, and you can hunt for it in your own projects. Look for anything being constructed, compiled or parsed inside a loop whose length is set by configuration rather than by your data: partitions, batch size, block size, number of tasks. Then ask whether its inputs ever change from one pass to the next. If they don’t, lift it out to the smallest scope where it’s still correct (per expression, per task, per stage, per executor), and write a test proving that the things meant to stay separate, like each task’s file list, the compression level or the regex pattern, didn’t get shared by accident. When you benchmark, push the multiplier to its extreme, but report the default as well. Most of these gains vanish at 200 partitions, and being upfront about that is what makes the 10,000-partition number believable.

What’s Next

1.1.0 branched off main on September 25. Seven more of my PRs landed after that and will ship in the next release. Two of them make Comet’s native metrics add up correctly in Spark’s task metrics, and one keeps shuffle reads native when adaptive query execution (AQE) reuses parts of the original plan. The native Delta scan from my last post, #5365, is still in review.

If you run Spark and haven’t given Comet a try, 1.1.0 is a good place to start. The installation guide covers the jars and settings you need.