Uh oh!
There was an error while loading. Please reload this page.
Core, Spark: Add JMH benchmarks for Variants - #15629
Conversation
7c4f806 to
2be00b9Compare@rashworld-max still a WiP I'm afraid. Need to know I'm measuring the right thing. Also I can't tell from your profile whether or not you are a human. |
70c69f8 to
25c6f29Compare| class ParquetVariantUtil { | ||
| @VisibleForTesting | ||
| public final class ParquetVariantUtil { |
There was a problem hiding this comment.
Is it possible to relocate the tests rather than expose this? We can do this, but generally prefer not to if we can avoid it.
There was a problem hiding this comment.
happy to do that...I have done it in the parquet PR
steveloughran
commented
Apr 14, 2026
I think this is ready for review. I've got the initial results and it's good for PRs like #3477 to be able to before/after benchmarks. More stuff can go in later; I've outlined them in my report. Equality deletes would be a fun one |
| */ | ||
| private long materializeNonEmpty(String operation, Dataset<?> ds) { | ||
| LOG.info("{} table={}", operation, tableType); | ||
| final long count = ds.count(); |
There was a problem hiding this comment.
Spark doesn't need to evaluate projection to count records.
There was a problem hiding this comment.
I needed something to do the entire compute and count() worked. Otherwise it's evaluate every row and feed to a black hole. What do you prefer?
There was a problem hiding this comment.
Spark can count records without evaluating projection so it's not really testing the projection here.
There was a problem hiding this comment.
it seemed to work, but I will get and discard each record instead
There was a problem hiding this comment.
@manuzhang fyi, now using the same sequence as the IcebergSourceBenchmark superclass, with the retention of the count for use in assertions
final long count = ds.queryExecution().toRdd().toJavaRDD().count();
blackhole.consume(count);
| "variant_get(nested, '$.varcategory', 'int')"; | ||
| /** Get the ID field from inside the variant: {@value}. */ | ||
| private static final String VARIANT_GET_NESTED_ID = "variant_get(nested, '$.varid', 'int')"; |
There was a problem hiding this comment.
should this be int64 as in the comments above?
There was a problem hiding this comment.
moved to bigint. Not doing enough with that type to measure any difference between that and int processing in the workflow.
steveloughran
commented
May 8, 2026
switching to draft again as I'm reworking the benchmark to show variant rowgroup filtering of shredded variants works (#15510), with changes including
Needs an extra PR in iceberg from qlong and a snaspshot of spark 4.1 with his changes for spark to pass variant_get down So, surprisingly complex. If I can show the chain works then it's time to start with the feature merges in spark and then here. This branch will merge without direct dependency, it's just a key goal of the spark benchmark is "show pushdown working". It's not ready to merge unless it can do that |
20d7bda to
2690e5cCompareFixesapache#15628 Core: benchmark of variant creation and ser/deser costs. Separate benchmarks for * building * serializing a prebuilt object * deserializing Variables are: - fields: [1000, 10000] - depth: [shallow, nested] - percentage of fields shed [0, 33, 67, 100] Note: the current benchmark does NOT for the JVM, as it allows for fast iterative development. A final merge should switch to fork(1). Spark: Full test of predicate pushdown of variants - avro - parquet unshredded - parquet shredded For this to return useful numbers, requires PRs for - Passing down variant_get between spark and iceberg - ParquetRowGroupFilter to filter on shredded variants. contains Add some more benchmarks Change-Id: I4231280f08cf63db5960ecb79301ae9458b35272
53dc868 to
34f347cCompare* file skipping can be observed as out of range equals/element tests are skipped completely * varcat filters still very slow Changes - blackhole consumption of rows - wiring up to ParquetMetricsRowGroupFilter.resetShreddedMetricsCounter() shows the filtering is happening on shredded files - cutting back on category count and attempting to change structure of file Note: this is the benchmark branch, and has had the assertions on counters and counter reset cut; all the other wiring up for the assertions is present. Change-Id: I686be12d51b13d2048b631b8cf198651012cc474
Allows for assertions in tests and in benchmarks that rowgroup skipping is taking place. Needed as there's not much tangible speedup, yet Change-Id: I8c03eb33d2d3d8a2139c347e6a72a7284e627f62
...which shows the configuration changes needed for data to be saved to multiple rowgroups. File size shrunk; increasing iterations of runs. Change-Id: Icb622958a068ed67de5bb895d88a9aa1713d2b11
- increasing category count reduces # of matches on the single category, so amplifying shredding filtering advantage - and varcat select with a range > 0 and < 1. That's the same as the = 1 and `in (1)`, selections, but with two scans of the values. Change-Id: Ib2b139697e235cb4674503784c6c909a5c460d1a
- fork=1, repetitions=5 - review feedback: use bigint in variant_get() on long values. Remove all commented out code referring to the predicate pushdown; everything is lined up waiting for that PR.
steveloughran
commented
May 28, 2026
@nssalian what do you think? |
steveloughran
commented
Jun 4, 2026
@rdblue this pr is ready for review in that "the data it generates, the layout and the queries are sufficient to show tangible speedup with qlong's file skipping and my predicate pushdown changes" |
steveloughran
commented
Jun 24, 2026
still ready for review |
steveloughran
commented
Jun 26, 2026
| } | ||
| /** Log the parquet metrics invocations. */ | ||
| private void logRowGroupFiltering() {} |
There was a problem hiding this comment.
logRowGroupFiltering(), resetMetricsCounter(), and expectFilteringOfShreddedFields() are all empty here but their javadoc describes behavior that doesn't exist, and expectPushdownWhenShredded is plumbed through select/selectEmptyResults just to gate the no-op expectFilteringOfShreddedFields(). Could these either be implemented or removed (along with the expectPushdownWhenShredded param)? Leaving them empty is a bit misleading to the next reader.
There was a problem hiding this comment.
ok, I will cut and when pushdown is added, they can be reinstated
| * <p>TODO: failing in spark. Needs investigation. | ||
| */ | ||
| // @Benchmark | ||
| public void joinArrElt4OnVarcat(Blackhole blackhole) { |
There was a problem hiding this comment.
Could we drop the commented-out joinArrElt4OnVarcat and open an issue for the failing self-join instead? A linked issue tracks it better than a // @Benchmark + TODO.
| /** Number of distinct category values. */ | ||
| private static final int NUM_CATEGORIES = 10; | ||
| public static final String COL_ID = "id"; |
There was a problem hiding this comment.
Minor: COL_ID and EQUALITY_CHECK are public while COL_NESTED/COL_ARR right below are private, and the sibling benchmarks keep their constants private. Same for TableType (and Depth/DEEP_NEST_DEPTH/ITERATIONS in the core benchmark). Worth making these private to match.
There was a problem hiding this comment.
I'll make them all private,inevitably just ide refactoring setup
steveloughran
left a comment
There was a problem hiding this comment.
@nssalian I will address your comments. Now that spark 4.2 is out we can make a bit more progress with variants, though I will have to move this pr onto being onto the spark/v4.2/ path as soon as it is added.
w.r.t the empty assertions, I will cut. for the curious here are the full ones
/** Log the parquet metrics invocations. */
private void logRowGroupFiltering() {
final long scans = ParquetMetricsRowGroupFilter.variantPredicatesShreddedMetricsEvaluated();
final long skipped = ParquetMetricsRowGroupFilter.variantPredicatesShreddedSkipped();
LOG.info("Scanned {} shredded metrics, skipped {} row groups", scans, skipped);
}
/** Reset the metrics counter. */
private void resetMetricsCounter() {
ParquetMetricsRowGroupFilter.resetShreddedMetricsCounters();
}
/**
* On shredded tables, assert that shredded field were scanned for filtering operations. Log the
* count at info.
*/
private void expectFilteringOfShreddedFields() {
if (isShredded()) {
final long scans = ParquetMetricsRowGroupFilter.variantPredicatesShreddedMetricsEvaluated();
final long skipped = ParquetMetricsRowGroupFilter.variantPredicatesShreddedSkipped();
LOG.info("Scanned {} shredded metrics, skipped {} row groups", scans, skipped);
assertThat(scans)
.describedAs("Number of times rowgroup metrics of shredded fields were scanned")
.isGreaterThan(0);
assertThat(skipped).describedAs("rowgroups skipped").isGreaterThan(0);
}
}
+ review comments. Change-Id: I44a5d6815cb8113bdf8ada8645089f859807f7bf
This pull request has been marked as stale due to 30 days of inactivity. It will be closed in 1 week if no further activity occurs. If you think that’s incorrect or this pull request requires a review, please simply write any comment. If closed, you can revive the PR at any time and @mention a reviewer or discuss it on the dev@iceberg.apache.org list. Thank you for your contributions. |
steveloughran
commented
Aug 21, 2026
@manuzhang thanks for marking as not stale. I am now ex cloudera, laptopless and mostly restricting my coding to learning golang for the fun of it right now. Maybe I'll get set up for java dev next week. |
nssalian
commented
Aug 23, 2026
@steveloughran it might be hard to get this in this release but I can try doing this after if you don't get a chance. |
steveloughran
commented
Aug 24, 2026
I'm not set up to do any java dev right now, so lets target the next release. The big source of merge pain will be the move to spark 4.2; when that goes in the pr will need the spark oart moved, but it will be straightforward |



Fixes#15628
core:VariantSerializationBenchmark
Separate benchmarks for
Variables are:
spark-4.1:IcebergSourceVariantReadBenchmark
Generate Avro, unshedded Parquet and shedded Parquet tables with the same variant data and then compare performance for basic filter and project operations against the normal columns and the variant fields.
Key findings:
I'm not reaching any conclusion why this is the case. I am looking at improving the performance of reconstructing string fields in parquet-java as those benchmark show needless byte-string-byte conversion. For the iceberg benchmark and layers below, I think knowing where issues like is enough of a change.
Writeup
See https://steveloughran.github.io/benchmarking-variants/ for the writeup and the interactive benchmark results of Iceberg and Parquet benchmarks.