From 5237fdd0f2a99e8e0fb63fb566d6bc94ce78be7d Mon Sep 17 00:00:00 2001 From: Anoop Date: Sat, 26 Sep 2026 00:51:09 +0530 Subject: [PATCH 1/4] fix(reader): keep column-free conjuncts in the row filter (#521) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Some queries return rows that do not match their `WHERE` clause. ```sql SELECT s FROM t WHERE NOT (s = s) ``` `s = s` is NULL wherever `s` is NULL, so `NOT (s = s)` is never true and this should return nothing. On `main` it returns every row whose `s` is NULL. ## Why `build_row_filter` breaks a conjunction into separate predicates (`split_conjunction`), then keeps only the conjuncts that reference a column: ```rust if required_indices_into_file_schema.is_empty() { return Ok(None); // conjunct dropped } ``` Simplification produces conjuncts that reference no column. `NOT (s = s)` becomes `s IS NULL AND NULL` — a column conjunct, and a literal `Boolean(NULL)` that reads nothing. The literal hit that branch and was dropped. Dropping it would be harmless if something else still applied it, but nothing does: DataFusion removes the `FilterExec` once it believes the predicate is fully pushed down, so the scan is the only place the predicate is applied. A dropped conjunct makes the filter strictly weaker than the one the user wrote, and rows that cannot match come back. ## The fix Keep the conjunct. A column-free candidate builds with an empty projection mask, and the cached path evaluates it against a batch carrying only the row count the selection implies — `RecordBatch` needs that stated explicitly, since there is no array present to imply it. ## Reproducing By hand: write a parquet file with a nullable `Utf8` column `s` where some rows are NULL, register it through `LiquidCacheLocalBuilder`, and run `SELECT s FROM t1 WHERE NOT (s = s)`. Expect no rows; on `main` every NULL-`s` row comes back. With the tests in this PR applied to `main` and the `row_filter.rs` change reverted: ``` $ cargo test -p liquid-cache-datafusion-local constant_conjunct test constant_null_conjunct_is_still_applied ... FAILED test three_way_partition_reconstructs_the_scan ... FAILED ---- constant_null_conjunct_is_still_applied ---- left: ["", "", ...] # 4000 rows returned right: [] # none expected ---- three_way_partition_reconstructs_the_scan ---- left: 14664 # rows across the three partitions right: 12000 # rows in the table test result: FAILED. 0 passed; 2 failed ``` ## The tests `src/datafusion-local/src/tests/constant_conjunct.rs`, both cold and warm cache: - **`constant_null_conjunct_is_still_applied`** — the direct case above. The fixture is 12000 rows with every third `s` NULL, giving the 4000. - **`three_way_partition_reconstructs_the_scan`** — the general property, and the more useful regression guard: for any predicate `P`, `WHERE P`, `WHERE NOT P` and `WHERE P IS NULL` must together return each row exactly once. Here rows whose `P` was NULL came back from both `WHERE NOT P` and `WHERE P IS NULL`. Found by a randomized query fuzzer using that three-way partition as its oracle. --- .../src/tests/constant_conjunct.rs | 131 ++++++++++++++++++ src/datafusion-local/src/tests/mod.rs | 1 + src/datafusion/src/cache/mod.rs | 13 +- .../src/reader/plantime/row_filter.rs | 9 +- 4 files changed, 149 insertions(+), 5 deletions(-) create mode 100644 src/datafusion-local/src/tests/constant_conjunct.rs diff --git a/src/datafusion-local/src/tests/constant_conjunct.rs b/src/datafusion-local/src/tests/constant_conjunct.rs new file mode 100644 index 00000000..6641372b --- /dev/null +++ b/src/datafusion-local/src/tests/constant_conjunct.rs @@ -0,0 +1,131 @@ +//! A pushed-down conjunct that references no column. +//! +//! Expression simplification turns `NOT (s = s)` into `s IS NULL AND NULL`: a +//! column conjunct and a literal `Boolean(NULL)` one. The literal reads no +//! column, and `build_row_filter` used to drop such a conjunct. Since DataFusion +//! removes the `FilterExec` when it pushes a predicate down, the scan is the only +//! place the predicate is applied, so dropping a conjunct *widens* the filter and +//! rows that cannot match come back. +//! +//! The visible symptom is a three-way partition that does not reconstruct the +//! scan: for a predicate `P`, `WHERE P`, `WHERE NOT P` and `WHERE P IS NULL` must +//! together return each row exactly once. Rows whose `P` is NULL were returned by +//! both `WHERE NOT P` and `WHERE P IS NULL`. + +use std::path::Path; +use std::sync::Arc; + +use arrow::array::{Array, Float64Array, Int64Array, RecordBatch, StringArray}; +use arrow_schema::{DataType, Field, Schema}; +use datafusion::prelude::{ParquetReadOptions, SessionConfig, SessionContext}; +use parquet::arrow::ArrowWriter; +use tempfile::TempDir; + +use crate::LiquidCacheLocalBuilder; + +/// The predicate under test. `s = s` is NULL wherever `s` is NULL, so a row with +/// a NULL `s`, an `f` above the threshold and an `id` outside the range makes the +/// whole predicate NULL. +const P: &str = "((t1.f <= 333.0 OR t1.s = t1.s) OR t1.id BETWEEN 1 AND 7)"; + +/// 12000 rows, more than the 8192-row default cache batch, so the scan spans more +/// than one cached batch. Every third row has a NULL `s`; `f` cycles 1..1000, so +/// two thirds of those NULL rows sit above the 333.0 threshold. +fn write_t1(path: &Path) { + let rows = 12000i64; + let schema = Arc::new(Schema::new(vec![ + Field::new("id", DataType::Int64, false), + Field::new("f", DataType::Float64, true), + Field::new("s", DataType::Utf8, true), + ])); + let id: Int64Array = (1..=rows).collect::>().into(); + let f: Float64Array = (1..=rows) + .map(|i| Some((i % 1000) as f64)) + .collect::>() + .into(); + let s: StringArray = (1..=rows) + .map(|i| (i % 3 != 0).then(|| format!("str{}", i % 17))) + .collect::>() + .into(); + let batch = + RecordBatch::try_new(schema.clone(), vec![Arc::new(id), Arc::new(f), Arc::new(s)]).unwrap(); + let file = std::fs::File::create(path).unwrap(); + let mut writer = ArrowWriter::try_new(file, schema, None).unwrap(); + writer.write(&batch).unwrap(); + writer.close().unwrap(); +} + +async fn liquid_ctx(dir: &Path) -> SessionContext { + let parquet = dir.join("t1.parquet"); + write_t1(&parquet); + let cache_dir = dir.join("cache"); + std::fs::create_dir_all(&cache_dir).unwrap(); + let (ctx, _cache) = LiquidCacheLocalBuilder::new() + .with_cache_dir(cache_dir) + .build(SessionConfig::new()) + .await + .unwrap(); + ctx.register_parquet( + "t1", + parquet.to_str().unwrap(), + ParquetReadOptions::default(), + ) + .await + .unwrap(); + ctx +} + +/// The `s` column as a sorted multiset, NULL rendered as ``. `s` is +/// projected as a string view, hence the cast. +async fn s_values(ctx: &SessionContext, sql: &str) -> Vec { + let batches = ctx.sql(sql).await.unwrap().collect().await.unwrap(); + let mut out = Vec::new(); + for batch in batches { + let column = arrow::compute::cast(batch.column(0), &DataType::Utf8).unwrap(); + let column = column.as_any().downcast_ref::().unwrap(); + for row in 0..column.len() { + out.push(match column.is_null(row) { + true => "".to_string(), + false => column.value(row).to_string(), + }); + } + } + out.sort(); + out +} + +/// `NOT (s = s)` is NULL where `s` is NULL and FALSE everywhere else, so it +/// matches no row. Simplification leaves it as `s IS NULL AND NULL`, and dropping +/// the literal conjunct returns every row with a NULL `s`. +#[tokio::test] +async fn constant_null_conjunct_is_still_applied() { + let dir = TempDir::new().unwrap(); + let ctx = liquid_ctx(dir.path()).await; + + // Cold reads through the source and fills the cache; warm is served from it, + // a separate evaluation path. + for pass in ["cold", "warm"] { + let rows = s_values(&ctx, "SELECT t1.s FROM t1 WHERE NOT (t1.s = t1.s)").await; + assert_eq!(rows, Vec::::new(), "{pass}"); + } +} + +/// `WHERE P`, `WHERE NOT P` and `WHERE P IS NULL` partition the table: together +/// they must return exactly the rows of the unfiltered scan, each once. +#[tokio::test] +async fn three_way_partition_reconstructs_the_scan() { + let dir = TempDir::new().unwrap(); + let ctx = liquid_ctx(dir.path()).await; + + for pass in ["cold", "warm"] { + let unfiltered = s_values(&ctx, "SELECT t1.s FROM t1").await; + + let mut partitioned = s_values(&ctx, &format!("SELECT t1.s FROM t1 WHERE {P}")).await; + partitioned.extend(s_values(&ctx, &format!("SELECT t1.s FROM t1 WHERE NOT {P}")).await); + partitioned.extend(s_values(&ctx, &format!("SELECT t1.s FROM t1 WHERE {P} IS NULL")).await); + partitioned.sort(); + + assert_eq!(partitioned.len(), unfiltered.len(), "{pass}"); + assert_eq!(partitioned, unfiltered, "{pass}"); + } +} diff --git a/src/datafusion-local/src/tests/mod.rs b/src/datafusion-local/src/tests/mod.rs index b5486a27..8bc6ebf1 100644 --- a/src/datafusion-local/src/tests/mod.rs +++ b/src/datafusion-local/src/tests/mod.rs @@ -19,6 +19,7 @@ use datafusion::{ }; use crate::LiquidCacheLocalBuilder; +mod constant_conjunct; mod date_optimizer; mod filter_limit; mod nested_filter; diff --git a/src/datafusion/src/cache/mod.rs b/src/datafusion/src/cache/mod.rs index fed837f6..d51bc611 100644 --- a/src/datafusion/src/cache/mod.rs +++ b/src/datafusion/src/cache/mod.rs @@ -5,7 +5,7 @@ use crate::io::ParquetCacheMetadata; use crate::reader::{LiquidPredicate, extract_multi_column_or}; use crate::sync::{Mutex, RwLock}; use ahash::AHashMap; -use arrow::array::{BooleanArray, RecordBatch}; +use arrow::array::{BooleanArray, RecordBatch, RecordBatchOptions}; use arrow::buffer::BooleanBuffer; use arrow_schema::{ArrowError, Field, Schema, SchemaRef}; use datafusion::common::tree_node::{Transformed, TreeNode}; @@ -242,6 +242,17 @@ impl CachedRowGroup { } } } + // A conjunct that reads no column still has to be evaluated, over a batch + // that carries only the row count the selection implies. + if column_ids.is_empty() { + let options = + RecordBatchOptions::new().with_row_count(Some(selection.count_set_bits())); + let record_batch = + RecordBatch::try_new_with_options(Arc::new(Schema::empty()), Vec::new(), &options) + .ok()?; + return Some(predicate.evaluate(record_batch)); + } + // Otherwise, we need to first convert the data into arrow arrays. let mut arrays = Vec::new(); let mut fields = Vec::new(); diff --git a/src/datafusion/src/reader/plantime/row_filter.rs b/src/datafusion/src/reader/plantime/row_filter.rs index 54c778ce..5d17d410 100644 --- a/src/datafusion/src/reader/plantime/row_filter.rs +++ b/src/datafusion/src/reader/plantime/row_filter.rs @@ -275,10 +275,11 @@ impl FilterCandidateBuilder { return Ok(None); }; - if required_indices_into_file_schema.is_empty() { - return Ok(None); - } - + // A conjunct that references no column - a literal `NULL` or `false` left + // behind by expression simplification, for instance - is still a conjunct. + // Dropping it here widens the filter, because by the time this runs + // DataFusion has removed the `FilterExec` on the assumption the predicate + // was fully pushed down, so the scan is the only place it is applied. let projected_file_schema = Arc::new( self.file_schema .project(&required_indices_into_file_schema)?, From 8d146620aef29e9b379a80f6d848b7dbeae9c59a Mon Sep 17 00:00:00 2001 From: Anoop Narang Date: Mon, 28 Sep 2026 09:46:37 +0530 Subject: [PATCH 2/4] docs: state what does not change, and check the remote names MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The file pinned upstream at a commit id. That was true the hour it was written and wrong three days later when our own PR merged upstream, and it would have gone on rotting on every upstream merge. A file whose job is telling agents what is true cannot hold facts that expire. It now says `main` contains upstream's history in full, which stays true, and gives the commands that answer where either side actually is. It also said "the remote is `origin` in a fresh clone", which reads as "origin is ours". In this working copy `origin` is *upstream* and the fork is a second remote named `fork` — the reverse. A session following this file wrote `git merge origin/main` meaning merge upstream, correctly for this clone, and it read as merging our own main into itself. There is now a Remotes section saying to run `git remote -v` first and warning that the names are not what you would guess, and every command names repositories by URL. Adds a Syncing section, because the omission is what sent that session looking. It says to merge upstream into a branch and raise a PR rather than using GitHub's "Sync fork" button, which offers to discard our commits when the merge is not a fast-forward, and it describes the conflict our own upstreamed changes cause on the way back: the content matches but the commit does not, so both sides appear to have edited the same lines. Resolve by keeping ours, and sync promptly, because alone that conflict is obvious and bundled with real upstream work it is not. Two more expiring claims removed. The shuttle note counted the failures it expects; it now says which tests fail and why. The toolchain note deferred to a `rust-toolchain.toml` that does not exist in this repo -- written from what such a repo usually has rather than from this one -- and now says there is none, which is the reason the pin is needed. --- CLAUDE.md | 73 +++++++++++++++++++++++++++++++++++++++++++++---------- 1 file changed, 60 insertions(+), 13 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index be5bc539..fa028016 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -1,20 +1,44 @@ # Working in this fork `hotdata-dev/liquid-cache` is a fork of `datafusion-contrib/liquid-cache`. -`main` is upstream's `0033b15` with our patches on top — upstream's history is -fully contained, so `git merge-base main ` resolves. +`main` contains upstream's history in full plus our patches, so +`git merge-base main ` resolves. + +Nothing here states where upstream currently is, or where we are: both move. +Ask git instead. + +``` +git fetch https://github.com/datafusion-contrib/liquid-cache main +git log --oneline FETCH_HEAD..main # ours that upstream does not have +git log --oneline main..FETCH_HEAD # upstream's that we do not have +``` `AGENTS.md` and `README.md` are upstream's and describe the project itself. This file is ours and describes only what differs here. Nothing in this repo should edit an upstream-owned file to record a fork convention: that conflicts on every sync. Add a file upstream does not have instead. +## Remotes + +**Check before using a remote name. They vary by clone and they are not what +you would guess** — in at least one working copy `origin` is *upstream* and the +fork is a second remote named `fork`, which is the reverse of the usual +arrangement. + +``` +git remote -v +``` + +Commands below name repositories by URL rather than by remote, so they are +correct in any clone. Do the same when writing instructions for anyone else: a +bare `origin/main` is ambiguous here and has already been misread as the +opposite of what it meant. + ## Branches -- Work off `main`. Fetch first — a local `main` goes stale with no signal, and - branching off a stale one silently drops everything merged since. The remote - is `origin` in a fresh clone; if you cloned upstream and added this fork as a - second remote, use that name instead. +- Work off `main`, and fetch before branching — a local `main` goes stale with + no signal, and branching off a stale one silently drops everything merged + since. - **Upstream PRs branch from upstream, not from `main`**, and are named `upstream/`. A branch cut from `main` carries our whole patch stack into the PR diff. @@ -26,8 +50,9 @@ on every sync. Add a file upstream does not have instead. Before pushing, this must list only the commits you wrote — anything else is a fork patch that would land in the upstream diff. Re-fetch upstream on the - line above it: `FETCH_HEAD` holds whatever the last fetch wrote, so after a - `git fetch origin` it is *our* `main` and the check hides every fork patch. + line above it: `FETCH_HEAD` holds whatever the last fetch wrote, so after + fetching any other remote it is no longer upstream and the check hides every + fork patch. ``` git fetch https://github.com/datafusion-contrib/liquid-cache main @@ -46,16 +71,38 @@ on every sync. Add a file upstream does not have instead. have (`DiskResidue`, `reclaim_orphaned_disk`, `settle`) and are not upstreamable at all. +## Syncing from upstream + +Merge upstream into a branch off `main` and raise a PR; do not use GitHub's +"Sync fork" button, which offers to discard our commits when the merge is not +a fast-forward. + +``` +git fetch https://github.com/datafusion-contrib/liquid-cache main +git checkout -b sync/upstream- main +git merge FETCH_HEAD +``` + +A change we contributed upstream comes back as their squash of it. The content +matches but the commit does not, so the merge conflicts where both sides +touched the same lines — typically a module list that each side appended to. +Resolve by keeping ours, which already contains the change. Sync promptly +rather than letting such a conflict wait: alone it is obvious, bundled with +real upstream work later it is not. + ## Building and testing -- **`cargo +1.96.0`.** The dependency tree needs 1.95+ (`vortex-*`, `sysinfo`) - and DataFusion 55 needs 1.94. A bare `cargo` on an older default fails - resolution with a wall of `requires rustc 1.9x` lines. +- **Pin the toolchain: `cargo +1.96.0 ...`.** There is no + `rust-toolchain.toml`, so a bare `cargo` uses whatever default is installed, + and an older one fails resolution with a wall of `requires rustc 1.9x` lines + naming `vortex-*` and `sysinfo`. Those lines state the minimum each crate + wants; use a toolchain at least that new. - **Shuttle tests need a filter**: `cargo +1.96.0 test -p liquid-cache --features shuttle --lib shuttle_`. The feature swaps `crate::sync` to shuttle primitives for the whole test build, so running it unfiltered fails - ~49 unrelated tests with "Are you accessing a Shuttle primitive outside of a - Shuttle test?". That is by design, not a regression. + every test that touches a lock outside a shuttle runner, with "Are you + accessing a Shuttle primitive outside of a Shuttle test?". That is by design, + not a regression. - **`dev-tools` needs `dev/dev-tools/assets/tailwind.css`**, which CI generates and the repo does not carry. To run its tests locally, create a placeholder and delete it before committing. Do not habitually pass From 833b8da4a6799dff9c587559db01286db815c9f4 Mon Sep 17 00:00:00 2001 From: Anoop Narang Date: Mon, 28 Sep 2026 10:12:04 +0530 Subject: [PATCH 3/4] docs: fix the sync recipe's own stale-main hazard, and narrow "keep ours" MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review found the Syncing section walking into the trap the Branches section warns about. Its only fetch named upstream, so local `main` was never refreshed, and the sync branch was cut from it — a stale one makes the sync PR revert fork commits merged since, with nothing to say so. It now fetches this fork first and branches from what came back, and uses each `FETCH_HEAD` immediately, since it holds only the last fetch. My own sync branch escaped this only because I had refreshed local `main` an hour earlier for unrelated reasons. "Resolve by keeping ours" was too broad. It is right for the returned squash of our own change, but a conflict hunk can hold an unrelated upstream edit on the same lines — a module appended next to the returned one — and taking the whole hunk from our side drops it silently. Now: keep ours for the returned change, keep any other upstream edit in the same hunk, and read the hunk rather than resolving by rule. Also records that a sync PR must be merged and not squashed. A squash gives the result one parent, so upstream's history never enters `main`'s ancestry: the merge-base does not move and the same upstream commits stay missing, with no error. Simulated both onto `main` — squashed, the merge-base stayed at the old upstream commit and the sync was undone; merged, it advanced to upstream's tip. The section now says to verify `git merge-base main FETCH_HEAD` afterwards. --- CLAUDE.md | 24 ++++++++++++++++++++---- 1 file changed, 20 insertions(+), 4 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index fa028016..773fa4d6 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -77,18 +77,34 @@ Merge upstream into a branch off `main` and raise a PR; do not use GitHub's "Sync fork" button, which offers to discard our commits when the merge is not a fast-forward. +Fetch *this fork* first and branch from what came back, not from local `main` +— the upstream fetch below never refreshes `main`, so a stale one would make +the sync PR revert fork commits merged since. Use each `FETCH_HEAD` +immediately: it holds only the last fetch. + ``` +git fetch https://github.com/hotdata-dev/liquid-cache main +git checkout -b sync/upstream- FETCH_HEAD + git fetch https://github.com/datafusion-contrib/liquid-cache main -git checkout -b sync/upstream- main git merge FETCH_HEAD ``` +**Merge this PR, do not squash it.** A squash gives the result a single parent, +so upstream's history never enters `main`'s ancestry: the merge-base does not +move, the same upstream commits stay missing, and nothing reports it. Verify +after merging that `git merge-base main FETCH_HEAD` is upstream's tip. + A change we contributed upstream comes back as their squash of it. The content matches but the commit does not, so the merge conflicts where both sides touched the same lines — typically a module list that each side appended to. -Resolve by keeping ours, which already contains the change. Sync promptly -rather than letting such a conflict wait: alone it is obvious, bundled with -real upstream work later it is not. +Keep ours *for the returned change*, and keep any other upstream edit in the +same hunk: upstream may have appended something of its own next to it, and +taking the whole hunk from our side drops that silently. Read the hunk rather +than resolving by rule. + +Sync promptly rather than letting such a conflict wait: alone it is obvious, +bundled with real upstream work later it is not. ## Building and testing From a4fd5e6b9333b074a0e6ace8794754feac6ce339 Mon Sep 17 00:00:00 2001 From: Anoop Narang Date: Mon, 28 Sep 2026 10:14:32 +0530 Subject: [PATCH 4/4] docs: make the post-sync check able to fail MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Review found the verify line checking neither side of what it claimed. `git merge-base main FETCH_HEAD` reads local `main`, which does not contain a merge made on GitHub, and `FETCH_HEAD` is whatever was fetched last — fetch the fork to refresh `main` and it compares the fork against itself. Run as written it returns our own tip and looks like a pass. It now captures upstream's tip before fetching the fork over it, and tests the property directly. Run against the current state, with this PR unmerged, it correctly reports not-ok; the old line reported our own main and told the reader nothing. Third time in this file that `FETCH_HEAD` holding only the last fetch has produced a wrong instruction. The sync recipe above now says so where the fetches are, rather than leaving each site to remember it. --- CLAUDE.md | 14 ++++++++++++-- 1 file changed, 12 insertions(+), 2 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index 773fa4d6..20cd2478 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -92,8 +92,18 @@ git merge FETCH_HEAD **Merge this PR, do not squash it.** A squash gives the result a single parent, so upstream's history never enters `main`'s ancestry: the merge-base does not -move, the same upstream commits stay missing, and nothing reports it. Verify -after merging that `git merge-base main FETCH_HEAD` is upstream's tip. +move, the same upstream commits stay missing, and nothing reports it. + +Verify afterwards. The merge happens on GitHub, so local `main` does not have +it and must not be what you check; and `FETCH_HEAD` holds only the last fetch, +so capture upstream before fetching the fork over it: + +``` +git fetch https://github.com/datafusion-contrib/liquid-cache main +upstream=$(git rev-parse FETCH_HEAD) +git fetch https://github.com/hotdata-dev/liquid-cache main +git merge-base --is-ancestor "$upstream" FETCH_HEAD && echo ok +``` A change we contributed upstream comes back as their squash of it. The content matches but the commit does not, so the merge conflicts where both sides