Skip to content

Upgrade to arrow 59 and DataFusion 55 - #146

Open
AdamGS wants to merge 2 commits into
apache:mainfrom
AdamGS:main
Open

AdamGS wants to merge 2 commits into
apache:mainfrom
AdamGS:main

Conversation

@AdamGS

@AdamGS AdamGS commented Sep 15, 2026

Copy link
Copy Markdown

This PR updates the Apache Arrow and DataFusion to their most recent releases - 59 and 55 (respectively). This change also required updating the rust-toolchain to one that satisfies their MSRV.

Also wanted to say thank you for this crate! Its been extremely useful for benchmarking our geospatial work!

Signed-off-by: Adam Gutglick <adam@spiraldb.com>
Signed-off-by: Adam Gutglick <adam@spiraldb.com>
@AdamGS

AdamGS commented Sep 15, 2026

Copy link
Copy Markdown
Author

I think test_zone_deterministic_parts_generation still fails, it should test that the data itself is identical right? I think DF changes something and the records within it change a bit because the write isn't sorted. Any opinions about fixing this? I think its either adding some sort to the write, or doing a full sort before the comparison.

@prantogg

Copy link
Copy Markdown
Contributor

@AdamGS thank you for this, and for the kind words!

I think its either adding some sort to the write, or doing a full sort before the comparison.

It's a little worse than ordering unfortunately, the rows themselves are different. DataFusion 55 folds the --part LIMIT into the window's sort like so:

DF 54:  SortExec(id)               <- FilterExec(..., fetch=1561) <- scan
DF 55:  SortExec: TopK(fetch=1561) <- FilterExec(...)             <- scan

So part 1 becomes the 1561 smallest ids instead of the first 1561 rows in scan order (11 of 1561 overlap with the reference). The zone parts are benchmark data that benchmark/answers/ depends on, so I think we need to keep the old selection rather than re-baseline. Materializing the LIMIT before the window step in generate_zone_parquet_single() does it:

let df = partition.apply_to_dataframe(df)?;
// Materialize the LIMIT before the window step. DataFusion otherwise folds it into a
// TopK over the ORDER BY id sort, which changes which rows land in this part.
let schema = Arc::clone(df.schema().inner());
let batches = df.collect().await?;
let df = ctx.read_table(Arc::new(datafusion::datasource::MemTable::try_new(
    schema,
    vec![batches],
)?))?;

With that, all six test_zone_* tests pass with the byte counts you already have in this PR 🙂. Everything else here looks good to me.

Do you happen to know if that TopK rewrite is intentional upstream? It also defeats early termination, so it might be worth an issue. I'll open a separate one here for making zone generation not depend on the optimizer at all, but definitely not for this PR!

@AdamGS

AdamGS commented Sep 16, 2026

Copy link
Copy Markdown
Author

There are a few DF issues/PRs around this sort of stuff, like:

Might be worth trying to backport them to 55.2 and try this upgrade again then.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants