From 377455993d06a2163ae508fa0bc81ec796b68561 Mon Sep 17 00:00:00 2001 From: Gabriel Date: Wed, 19 Aug 2026 12:40:41 +0200 Subject: [PATCH 1/5] test: add Iceberg integration fixture --- .gitattributes | 1 + Cargo.lock | 16 ++ crates/datafusion/Cargo.toml | 4 + crates/datafusion/src/lib.rs | 3 + crates/datafusion/src/test_utils/harness.rs | 160 +++++++++++++++++ crates/datafusion/src/test_utils/mod.rs | 3 + .../testdata/iceberg/taxi/README.md | 76 +++++++++ .../pickup_date=2024-01-08/data_0.parquet | 3 + .../pickup_date=2024-01-09/data_0.parquet | 3 + .../pickup_date=2024-01-10/data_0.parquet | 3 + .../pickup_date=2024-01-11/data_0.parquet | 3 + .../pickup_date=2024-01-12/data_0.parquet | 3 + .../pickup_date=2024-01-13/data_0.parquet | 3 + .../pickup_date=2024-01-14/data_0.parquet | 3 + ...-47e0-4c4b-9522-4a7c44d74036.metadata.json | 1 + ...9fdb82-eb66-7582-99a7-9f864b92a53f-m0.avro | Bin 0 -> 5141 bytes ...-019fdb82-eb66-7582-99a7-9f864b92a53f.avro | Bin 0 -> 1607 bytes .../iceberg/taxi/metadata/v1.metadata.json | 1 + crates/datafusion/tests/external_table.rs | 49 ++++++ crates/datafusion/tests/filter_pushdown.rs | 161 ++++++++++++++++++ crates/datafusion/tests/limit_pushdown.rs | 141 +++++++++++++++ .../datafusion/tests/projection_pushdown.rs | 140 +++++++++++++++ 22 files changed, 777 insertions(+) create mode 100644 .gitattributes create mode 100644 crates/datafusion/src/test_utils/harness.rs create mode 100644 crates/datafusion/src/test_utils/mod.rs create mode 100644 crates/datafusion/testdata/iceberg/taxi/README.md create mode 100644 crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-08/data_0.parquet create mode 100644 crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-09/data_0.parquet create mode 100644 crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-10/data_0.parquet create mode 100644 crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-11/data_0.parquet create mode 100644 crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-12/data_0.parquet create mode 100644 crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-13/data_0.parquet create mode 100644 crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-14/data_0.parquet create mode 100644 crates/datafusion/testdata/iceberg/taxi/metadata/00000-00a113a6-47e0-4c4b-9522-4a7c44d74036.metadata.json create mode 100644 crates/datafusion/testdata/iceberg/taxi/metadata/019fdb82-eb66-7582-99a7-9f864b92a53f-m0.avro create mode 100644 crates/datafusion/testdata/iceberg/taxi/metadata/snap-3167948105555765929-0-019fdb82-eb66-7582-99a7-9f864b92a53f.avro create mode 100644 crates/datafusion/testdata/iceberg/taxi/metadata/v1.metadata.json create mode 100644 crates/datafusion/tests/external_table.rs create mode 100644 crates/datafusion/tests/filter_pushdown.rs create mode 100644 crates/datafusion/tests/limit_pushdown.rs create mode 100644 crates/datafusion/tests/projection_pushdown.rs diff --git a/.gitattributes b/.gitattributes new file mode 100644 index 0000000..b57f642 --- /dev/null +++ b/.gitattributes @@ -0,0 +1 @@ +*.parquet filter=lfs diff=lfs merge=lfs -text diff --git a/Cargo.lock b/Cargo.lock index 60a2093..71887ed 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2112,14 +2112,18 @@ name = "datafusion-iceberg" version = "0.10.1" dependencies = [ "async-trait", + "bytes", "dashmap", "datafusion", "expect-test", "futures", "iceberg", + "insta", "parquet", + "serde", "tempfile", "tokio", + "typetag", "uuid", ] @@ -3464,6 +3468,18 @@ dependencies = [ "generic-array", ] +[[package]] +name = "insta" +version = "1.48.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "86f0f8fee8c926415c58d6ae43a08523a26faccb2323f5e6b644fe7dd4ef6b82" +dependencies = [ + "console", + "once_cell", + "similar", + "tempfile", +] + [[package]] name = "inventory" version = "0.3.24" diff --git a/crates/datafusion/Cargo.toml b/crates/datafusion/Cargo.toml index b4c3b92..6d59998 100644 --- a/crates/datafusion/Cargo.toml +++ b/crates/datafusion/Cargo.toml @@ -30,6 +30,7 @@ repository = { workspace = true } [dependencies] async-trait = { workspace = true } +bytes = "1" dashmap = "6" datafusion = { workspace = true } futures = "0.3" @@ -37,9 +38,12 @@ iceberg = { workspace = true } parquet = "59.2" tokio = { workspace = true } uuid = { version = "1.18", features = ["v7"] } +serde = { version = "1", features = ["derive"] } +typetag = "0.2" [dev-dependencies] expect-test = "1" +insta = "1" parquet = "59.2" tempfile = "3.18" diff --git a/crates/datafusion/src/lib.rs b/crates/datafusion/src/lib.rs index 4b0ea86..fd18f08 100644 --- a/crates/datafusion/src/lib.rs +++ b/crates/datafusion/src/lib.rs @@ -27,4 +27,7 @@ pub mod table; pub use table::table_provider_factory::IcebergTableProviderFactory; pub use table::*; +#[doc(hidden)] +pub mod test_utils; + pub(crate) mod task_writer; diff --git a/crates/datafusion/src/test_utils/harness.rs b/crates/datafusion/src/test_utils/harness.rs new file mode 100644 index 0000000..90e93c0 --- /dev/null +++ b/crates/datafusion/src/test_utils/harness.rs @@ -0,0 +1,160 @@ +use std::path::PathBuf; +use std::sync::Arc; + +use async_trait::async_trait; +use bytes::Bytes; +use datafusion::arrow::util::pretty::pretty_format_batches; +use datafusion::dataframe::DataFrame; +use datafusion::error::Result; +use datafusion::execution::SessionStateBuilder; +use datafusion::physical_plan::displayable; +use datafusion::prelude::SessionContext; +use futures::StreamExt; +use futures::stream::BoxStream; +use iceberg::io::{ + FileMetadata, FileRead, FileWrite, InputFile, LocalFsStorage, OutputFile, Storage, + StorageConfig, StorageFactory, +}; +use iceberg::{Error, ErrorKind, Result as IcebergResult}; +use serde::{Deserialize, Serialize}; + +use crate::IcebergTableProviderFactory; + +pub const FIXTURE_URI: &str = "s3://iceberg-test/warehouse/taxi"; +const WAREHOUSE_URI: &str = "s3://iceberg-test/warehouse/"; + +pub struct IcebergTestHarness { + ctx: SessionContext, +} + +impl IcebergTestHarness { + pub async fn new() -> Result { + let state = SessionStateBuilder::new() + .with_default_features() + .with_table_factory( + "ICEBERG".to_string(), + Arc::new(IcebergTableProviderFactory::new_with_storage_factory( + Arc::new(FixtureStorageFactory::default()), + )), + ) + .build(); + let ctx = SessionContext::new_with_state(state); + ctx.sql(&format!( + "CREATE EXTERNAL TABLE taxi STORED AS ICEBERG \ + LOCATION '{FIXTURE_URI}/metadata/v1.metadata.json'" + )) + .await? + .collect() + .await?; + Ok(Self { ctx }) + } + + pub async fn query(&self, sql: &str) -> Result<(String, String)> { + let dataframe: DataFrame = self.ctx.sql(sql).await?; + let plan = dataframe.create_physical_plan().await?; + let batches = dataframe.collect().await?; + + Ok(( + displayable(plan.as_ref()).indent(true).to_string(), + pretty_format_batches(&batches)?.to_string(), + )) + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct FixtureStorageFactory { + root: PathBuf, +} + +impl Default for FixtureStorageFactory { + fn default() -> Self { + Self { + root: PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("testdata/iceberg"), + } + } +} + +#[typetag::serde] +impl StorageFactory for FixtureStorageFactory { + fn build(&self, _config: &StorageConfig) -> IcebergResult> { + Ok(Arc::new(FixtureStorage { + root: self.root.clone(), + })) + } +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +struct FixtureStorage { + root: PathBuf, +} + +impl FixtureStorage { + fn local_path(&self, path: &str) -> IcebergResult { + let relative = path.strip_prefix(WAREHOUSE_URI).ok_or_else(|| { + Error::new( + ErrorKind::DataInvalid, + format!("unsupported fixture URI: {path}"), + ) + })?; + Ok(self.root.join(relative).display().to_string()) + } + + fn local(&self) -> LocalFsStorage { + LocalFsStorage::new() + } +} + +#[async_trait] +#[typetag::serde] +impl Storage for FixtureStorage { + async fn exists(&self, path: &str) -> IcebergResult { + self.local().exists(&self.local_path(path)?).await + } + + async fn metadata(&self, path: &str) -> IcebergResult { + self.local().metadata(&self.local_path(path)?).await + } + + async fn read(&self, path: &str) -> IcebergResult { + self.local().read(&self.local_path(path)?).await + } + + async fn reader(&self, path: &str) -> IcebergResult> { + self.local().reader(&self.local_path(path)?).await + } + + async fn write(&self, path: &str, bytes: Bytes) -> IcebergResult<()> { + self.local().write(&self.local_path(path)?, bytes).await + } + + async fn writer(&self, path: &str) -> IcebergResult> { + self.local().writer(&self.local_path(path)?).await + } + + async fn delete(&self, path: &str) -> IcebergResult<()> { + self.local().delete(&self.local_path(path)?).await + } + + async fn delete_prefix(&self, path: &str) -> IcebergResult<()> { + self.local().delete_prefix(&self.local_path(path)?).await + } + + async fn delete_stream( + &self, + paths: BoxStream<'static, String>, + ) -> IcebergResult<()> { + let mut paths = paths; + while let Some(path) = paths.next().await { + self.delete(&path).await?; + } + Ok(()) + } + + fn new_input(&self, path: &str) -> IcebergResult { + Ok(InputFile::new(Arc::new(self.clone()), path.to_string())) + } + + fn new_output(&self, path: &str) -> IcebergResult { + Ok(OutputFile::new(Arc::new(self.clone()), path.to_string())) + } +} diff --git a/crates/datafusion/src/test_utils/mod.rs b/crates/datafusion/src/test_utils/mod.rs new file mode 100644 index 0000000..d3fc355 --- /dev/null +++ b/crates/datafusion/src/test_utils/mod.rs @@ -0,0 +1,3 @@ +mod harness; + +pub use harness::*; diff --git a/crates/datafusion/testdata/iceberg/taxi/README.md b/crates/datafusion/testdata/iceberg/taxi/README.md new file mode 100644 index 0000000..b4fed84 --- /dev/null +++ b/crates/datafusion/testdata/iceberg/taxi/README.md @@ -0,0 +1,76 @@ +# NYC Yellow Taxi Iceberg fixture source data + +This directory is a complete, frozen Iceberg v2 table containing a 175,000-row +slice of the NYC Taxi & Limousine Commission Yellow Taxi Trip Record Data. +It includes data files, a manifest, a manifest list, and table metadata. +Each Parquet field carries the corresponding Iceberg field ID. + +The table uses the stable URI prefix `s3://iceberg-test/warehouse/taxi/` in +its metadata. Tests map that prefix to this checked-in directory, so the +metadata stays valid regardless of where the repository is cloned. + +## Selection + +- Source: `yellow_tripdata_2024-01.parquet`, published by the NYC TLC. +- Pickup dates: 2024-01-08 through 2024-01-14 (inclusive). +- 25,000 valid trips per pickup date, ordered by pickup timestamp and location + IDs, for a total of 175,000 rows. +- Excludes records with non-positive distance or total amount. +- Data is stored in seven Zstandard-compressed Parquet files, partitioned by + `pickup_date`. + +The original records contain pickup/drop-off timestamps and locations, +passenger count, distance, payment type, and fare components. This fixture +keeps those fields while using snake_case column names. + +## Provenance + +The NYC TLC publishes the monthly trip records as Parquet files and documents +the included fields at: + + + +## Source extraction + +The data slice was produced with DuckDB 1.4.2 before writing the committed +Iceberg metadata tree with Iceberg Rust 0.10.1: + +```sql +COPY ( + WITH sampled AS ( + SELECT + VendorID AS vendor_id, + tpep_pickup_datetime AS pickup_at, + tpep_dropoff_datetime AS dropoff_at, + passenger_count, + trip_distance, + PULocationID AS pickup_location_id, + DOLocationID AS dropoff_location_id, + payment_type, + fare_amount, + tip_amount, + tolls_amount, + total_amount, + CAST(tpep_pickup_datetime AS DATE) AS pickup_date, + row_number() OVER ( + PARTITION BY CAST(tpep_pickup_datetime AS DATE) + ORDER BY tpep_pickup_datetime, PULocationID, DOLocationID + ) AS row_number + FROM read_parquet( + 'https://d37ci6vzurychx.cloudfront.net/trip-data/yellow_tripdata_2024-01.parquet' + ) + WHERE tpep_pickup_datetime >= TIMESTAMP '2024-01-08 00:00:00' + AND tpep_pickup_datetime < TIMESTAMP '2024-01-15 00:00:00' + AND trip_distance > 0 + AND total_amount > 0 + ) + SELECT * EXCLUDE (row_number) + FROM sampled + WHERE row_number <= 25000 +) TO 'testdata/iceberg/taxi/data' ( + FORMAT parquet, + PARTITION_BY (pickup_date), + COMPRESSION zstd, + ROW_GROUP_SIZE 25000 +); +``` diff --git a/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-08/data_0.parquet b/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-08/data_0.parquet new file mode 100644 index 0000000..89b142f --- /dev/null +++ b/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-08/data_0.parquet @@ -0,0 +1,3 @@ +version https://git-lfs.github.com/spec/v1 +oid sha256:821513f8a2dd6a05df4de8fb7f6b52b4d8110ee2be6540b0fd67f74cc74eec65 +size 639826 diff --git a/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-09/data_0.parquet b/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-09/data_0.parquet new file mode 100644 index 0000000..338ba51 --- /dev/null +++ b/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-09/data_0.parquet @@ -0,0 +1,3 @@ +version https://git-lfs.github.com/spec/v1 +oid sha256:4ce6a4bde784ff55a6aa8efa6b1bdb8666e6292734698ce66ac65d2119e388e6 +size 620354 diff --git a/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-10/data_0.parquet b/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-10/data_0.parquet new file mode 100644 index 0000000..e82b673 --- /dev/null +++ b/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-10/data_0.parquet @@ -0,0 +1,3 @@ +version https://git-lfs.github.com/spec/v1 +oid sha256:dad448223c53f1327ad30fd440f429946006b3866cb16417b25a42cf20f219c9 +size 628861 diff --git a/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-11/data_0.parquet b/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-11/data_0.parquet new file mode 100644 index 0000000..3c12738 --- /dev/null +++ b/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-11/data_0.parquet @@ -0,0 +1,3 @@ +version https://git-lfs.github.com/spec/v1 +oid sha256:800c87ec48e55097974ee1233b2c77b186469811368081ef9acb1a5a5de0cc19 +size 627604 diff --git a/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-12/data_0.parquet b/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-12/data_0.parquet new file mode 100644 index 0000000..92100e2 --- /dev/null +++ b/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-12/data_0.parquet @@ -0,0 +1,3 @@ +version https://git-lfs.github.com/spec/v1 +oid sha256:8e1d7d78e847f376358b543b487d2ea0ea47363529ce3aa32bd6eb2db4fb64d0 +size 650045 diff --git a/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-13/data_0.parquet b/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-13/data_0.parquet new file mode 100644 index 0000000..765321d --- /dev/null +++ b/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-13/data_0.parquet @@ -0,0 +1,3 @@ +version https://git-lfs.github.com/spec/v1 +oid sha256:98bf67f74e8cc6626d1d8e30450aff86fdda081c8a1fa7961ac5123666dc4fd4 +size 660972 diff --git a/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-14/data_0.parquet b/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-14/data_0.parquet new file mode 100644 index 0000000..153558e --- /dev/null +++ b/crates/datafusion/testdata/iceberg/taxi/data/pickup_date=2024-01-14/data_0.parquet @@ -0,0 +1,3 @@ +version https://git-lfs.github.com/spec/v1 +oid sha256:28b7fb2209730ef2bc7c33626a2472366bbf1fe3846f38343b15bf7e5e079102 +size 652720 diff --git a/crates/datafusion/testdata/iceberg/taxi/metadata/00000-00a113a6-47e0-4c4b-9522-4a7c44d74036.metadata.json b/crates/datafusion/testdata/iceberg/taxi/metadata/00000-00a113a6-47e0-4c4b-9522-4a7c44d74036.metadata.json new file mode 100644 index 0000000..49a5e24 --- /dev/null +++ b/crates/datafusion/testdata/iceberg/taxi/metadata/00000-00a113a6-47e0-4c4b-9522-4a7c44d74036.metadata.json @@ -0,0 +1 @@ +{"format-version":2,"table-uuid":"019fdb82-eb63-7de3-b33f-f8063d4fe433","location":"s3://iceberg-test/warehouse/taxi","last-sequence-number":0,"last-updated-ms":1786094218083,"last-column-id":13,"schemas":[{"schema-id":0,"type":"struct","fields":[{"id":1,"name":"vendor_id","required":false,"type":"int"},{"id":2,"name":"pickup_at","required":false,"type":"timestamp"},{"id":3,"name":"dropoff_at","required":false,"type":"timestamp"},{"id":4,"name":"passenger_count","required":false,"type":"long"},{"id":5,"name":"trip_distance","required":false,"type":"double"},{"id":6,"name":"pickup_location_id","required":false,"type":"int"},{"id":7,"name":"dropoff_location_id","required":false,"type":"int"},{"id":8,"name":"payment_type","required":false,"type":"long"},{"id":9,"name":"fare_amount","required":false,"type":"double"},{"id":10,"name":"tip_amount","required":false,"type":"double"},{"id":11,"name":"tolls_amount","required":false,"type":"double"},{"id":12,"name":"total_amount","required":false,"type":"double"},{"id":13,"name":"pickup_date","required":false,"type":"date"}]}],"current-schema-id":0,"partition-specs":[{"spec-id":0,"fields":[{"source-id":13,"field-id":1000,"name":"pickup_date","transform":"identity"}]}],"default-spec-id":0,"last-partition-id":1000,"sort-orders":[{"order-id":0,"fields":[]}],"default-sort-order-id":0,"refs":{}} \ No newline at end of file diff --git a/crates/datafusion/testdata/iceberg/taxi/metadata/019fdb82-eb66-7582-99a7-9f864b92a53f-m0.avro b/crates/datafusion/testdata/iceberg/taxi/metadata/019fdb82-eb66-7582-99a7-9f864b92a53f-m0.avro new file mode 100644 index 0000000000000000000000000000000000000000..364cdaf654ed131d7252051420d12de6fd093d43 GIT binary patch literal 5141 zcmcIn&x<2P6wXADJ?J1Jp7zw>E!p&s%uEvS;9*aSv&xExWoRnh)idREclB0P&um-< z6x5q~^W<@nVOexRPyPuaA|mQV(3>E5^dKI5Roz{wPSTl4%kef52>Uhj|J z+Ir@ChL|9N(6q}l9DKTU=f1^h&H~&ep*3>6jujCchsFPF+wNE?O0YGuG7=o;nI9s7 zt&SxaN_j+C0&hr&Q^+~7&JM22M@V!}Fym08?LNzcBb*>e*S4QSCyd?*Xo!PrX&%Q< z74Lq08X#1^vPMu*Ofp%K3zi20iYY|g8Ua=XQe!&7X-Jt54?7macXPrpG!dbgW1WSh z!aC~|k*jpTDPlVfHA<74oEpXSLJGhP1j**qg# zL}Kjo|DCdLq=b=)e3@56Yb2Su6=B7=E6ECIDu8w!ZxzgHdO**{)b`Bcs-FyrqiHbA zRRSZW7Gs-dqhg9lenFIZ4q93|rnybWUfgJ1&2?4iY(?`iJtP5&Zx?rkx|YtCN(L@7 za6r9l*Eps^@Om7V8dg%Q2hGccyp4U5`s0bfhWr z=^xA%DX88|??dR_wb$u=I16h^)cgUug@jh_-7F$!nXg#hAOWr;30Q2Z<%jQ@9=`DK zU3aa+cg*pUqIBc^djsg-8?GS$Z%&uX1mF$V>fbXR7SmH$`^K=dgc}jRyDQIrZ;jsF z-DdGu!_2r?+`YA$cXx~`d6sRqdFgy_9p}5IDQASgpm}dS>$|4Bz=?i>1(F+qR#`Rj zt)l4CA?lfC#ja_=rb`@H4D+cPgpSdI`>p1|t1LfPWhsaM^- zpsFbxP8l@ILQs$T-m5PMY9m@oUEdr*5n-^&GJ3i&*Tw5+-{gbg2s62XhVtupX4TWD zTW(#;F`BDIx%8C>JTI3^pscSOYKK(&#{9CPVO}EZj+z3~xgf8E?Go)Z$v|uW&+gZc zkLfpGzVHru^4s_S-Tr91-TvX>7v~?p`^VpJwpy)EwzxOy^+mjPFwLMN2>|L|FZW~oOAaUZt zFCb2|Lcmw_132&v`~aL8Cvmc^r6?7N6xo}3GjDd@o0sA1y?Y1rREi@NO(CP#`xivd zav&ojArdksEy8J5zGsxP1e6XEmI59p4ANMU(fFbkkyAazh*ct`CF7F{VY8TwJja>0 z>bMm6#6+f|hI0*EPNim=2_C2%q5u=GYcj5H73>(mNN^2{h!w?JxH~Yy6}&BgM-cEL z3qfwG)cLudfy)>&y92gwU}+pf94Jn6H5K|sZ}qB9w}vc-q=HBkTibMc1`y7f(m1Kb zfOo)}t;UdoM#4Lmt#>NLS*>KfQP-T}n(vNS^O7pQJ3`HCTB}a2g)k<(NK=A}Ug?QS zNoV*s(~zM>*5|i9B|T?Y>GXz`N`+F!6-us5T;&agNP+SSvhtE|3~>Wogf010{78jA zX^|KbTBMqc$_gttFY~!WBVQ`cKn5Y|XydL;x|s$=wkI=kAH!uq_&@`r+u5a~D?fA&xpt2=0c-D3M1HU(P9S=W)zUO!RCy9_5)%Gb! zh1;pod5W4T2^+X(dF}g)%^QgrqQ(1RwJq%BrfIm=EA%?}dJc wdfmj%oTHM+)ho+tG=5%v|MKzO@4v_QE$hd(hvkul-{ylC_}F}ezgCUZKi`fLW&i*H literal 0 HcmV?d00001 diff --git a/crates/datafusion/testdata/iceberg/taxi/metadata/v1.metadata.json b/crates/datafusion/testdata/iceberg/taxi/metadata/v1.metadata.json new file mode 100644 index 0000000..9bddf41 --- /dev/null +++ b/crates/datafusion/testdata/iceberg/taxi/metadata/v1.metadata.json @@ -0,0 +1 @@ +{"format-version":2,"table-uuid":"019fdb82-eb63-7de3-b33f-f8063d4fe433","location":"s3://iceberg-test/warehouse/taxi","last-sequence-number":1,"last-updated-ms":1786094218149,"last-column-id":13,"schemas":[{"schema-id":0,"type":"struct","fields":[{"id":1,"name":"vendor_id","required":false,"type":"int"},{"id":2,"name":"pickup_at","required":false,"type":"timestamp"},{"id":3,"name":"dropoff_at","required":false,"type":"timestamp"},{"id":4,"name":"passenger_count","required":false,"type":"long"},{"id":5,"name":"trip_distance","required":false,"type":"double"},{"id":6,"name":"pickup_location_id","required":false,"type":"int"},{"id":7,"name":"dropoff_location_id","required":false,"type":"int"},{"id":8,"name":"payment_type","required":false,"type":"long"},{"id":9,"name":"fare_amount","required":false,"type":"double"},{"id":10,"name":"tip_amount","required":false,"type":"double"},{"id":11,"name":"tolls_amount","required":false,"type":"double"},{"id":12,"name":"total_amount","required":false,"type":"double"},{"id":13,"name":"pickup_date","required":false,"type":"date"}]}],"current-schema-id":0,"partition-specs":[{"spec-id":0,"fields":[{"source-id":13,"field-id":1000,"name":"pickup_date","transform":"identity"}]}],"default-spec-id":0,"last-partition-id":1000,"current-snapshot-id":3167948105555765929,"snapshot-log":[{"snapshot-id":3167948105555765929,"timestamp-ms":1786094218149}],"metadata-log":[{"metadata-file":"s3://iceberg-test/warehouse/taxi/metadata/00000-00a113a6-47e0-4c4b-9522-4a7c44d74036.metadata.json","timestamp-ms":1786094218083}],"sort-orders":[{"order-id":0,"fields":[]}],"default-sort-order-id":0,"refs":{"main":{"snapshot-id":3167948105555765929,"type":"branch"}},"snapshots":[{"snapshot-id":3167948105555765929,"sequence-number":1,"timestamp-ms":1786094218149,"manifest-list":"s3://iceberg-test/warehouse/taxi/metadata/snap-3167948105555765929-0-019fdb82-eb66-7582-99a7-9f864b92a53f.avro","summary":{"operation":"append","added-data-files":"7","total-records":"175000","changed-partition-count":"7","added-files-size":"4480382","total-delete-files":"0","total-position-deletes":"0","total-equality-deletes":"0","added-records":"175000","total-files-size":"4480382","total-data-files":"7"},"schema-id":0}]} \ No newline at end of file diff --git a/crates/datafusion/tests/external_table.rs b/crates/datafusion/tests/external_table.rs new file mode 100644 index 0000000..2f06e4e --- /dev/null +++ b/crates/datafusion/tests/external_table.rs @@ -0,0 +1,49 @@ +#[cfg(test)] +mod tests { + use datafusion::error::Result; + use datafusion_iceberg::test_utils::{FIXTURE_URI, IcebergTestHarness}; + + #[tokio::test] + async fn registers_the_fixture_with_the_iceberg_schema() -> Result<()> { + let harness = IcebergTestHarness::new().await?; + let (_, batches) = harness.query("DESCRIBE taxi").await?; + + insta::assert_snapshot!(batches, @r" + +---------------------+---------------+-------------+ + | column_name | data_type | is_nullable | + +---------------------+---------------+-------------+ + | vendor_id | Int32 | YES | + | pickup_at | Timestamp(µs) | YES | + | dropoff_at | Timestamp(µs) | YES | + | passenger_count | Int64 | YES | + | trip_distance | Float64 | YES | + | pickup_location_id | Int32 | YES | + | dropoff_location_id | Int32 | YES | + | payment_type | Int64 | YES | + | fare_amount | Float64 | YES | + | tip_amount | Float64 | YES | + | tolls_amount | Float64 | YES | + | total_amount | Float64 | YES | + | pickup_date | Date32 | YES | + +---------------------+---------------+-------------+ + "); + + Ok(()) + } + + #[tokio::test] + async fn rejects_schema_definitions_for_existing_iceberg_tables() -> Result<()> { + let harness = IcebergTestHarness::new().await?; + let error = harness + .query(&format!( + "CREATE EXTERNAL TABLE invalid (id INT) STORED AS ICEBERG \ + LOCATION '{FIXTURE_URI}/metadata/v1.metadata.json'" + )) + .await + .unwrap_err(); + + insta::assert_snapshot!(error.to_string(), @"This feature is not implemented: Currently we only support reading existing icebergs tables in external table command. To create new table, please use catalog provider."); + + Ok(()) + } +} diff --git a/crates/datafusion/tests/filter_pushdown.rs b/crates/datafusion/tests/filter_pushdown.rs new file mode 100644 index 0000000..d82f107 --- /dev/null +++ b/crates/datafusion/tests/filter_pushdown.rs @@ -0,0 +1,161 @@ +#[cfg(test)] +mod tests { + use datafusion::error::Result; + use datafusion_iceberg::test_utils::IcebergTestHarness; + + #[tokio::test] + async fn pushes_down_compound_predicates() -> Result<()> { + let harness = IcebergTestHarness::new().await?; + let (plan, batches) = harness + .query( + r"SELECT COUNT(*) AS trips + FROM taxi + WHERE pickup_date = DATE '2024-01-10' + AND payment_type IN (1, 2) + AND trip_distance >= 2.0", + ) + .await?; + + insta::assert_snapshot!(plan, @r" + ProjectionExec: expr=[count(Int64(1))@0 as trips] + AggregateExec: mode=Final, gby=[], aggr=[count(Int64(1))] + CoalescePartitionsExec + AggregateExec: mode=Partial, gby=[], aggr=[count(Int64(1))] + FilterExec: pickup_date@2 = 2024-01-10 AND (payment_type@1 = 1 OR payment_type@1 = 2) AND trip_distance@0 >= 2, projection=[] + RepartitionExec: partitioning=RoundRobinBatch(16), input_partitions=1 + IcebergTableScan projection:[trip_distance,payment_type,pickup_date] predicate:[((pickup_date = 2024-01-10) AND ((payment_type = 1) OR (payment_type = 2))) AND (trip_distance >= 2)] + "); + insta::assert_snapshot!(batches, @r" + +-------+ + | trips | + +-------+ + | 9891 | + +-------+ + "); + + Ok(()) + } + + #[tokio::test] + async fn pushes_down_disjunctions() -> Result<()> { + let harness = IcebergTestHarness::new().await?; + let (plan, batches) = harness + .query( + r"SELECT pickup_date, COUNT(*) AS trips + FROM taxi + WHERE pickup_date = DATE '2024-01-10' + OR pickup_date = DATE '2024-01-11' + GROUP BY pickup_date + ORDER BY pickup_date", + ) + .await?; + + insta::assert_snapshot!(plan, @r" + SortPreservingMergeExec: [pickup_date@0 ASC NULLS LAST] + SortExec: expr=[pickup_date@0 ASC NULLS LAST], preserve_partitioning=[true] + ProjectionExec: expr=[pickup_date@0 as pickup_date, count(Int64(1))@1 as trips] + AggregateExec: mode=FinalPartitioned, gby=[pickup_date@0 as pickup_date], aggr=[count(Int64(1))] + RepartitionExec: partitioning=Hash([pickup_date@0], 16), input_partitions=16 + AggregateExec: mode=Partial, gby=[pickup_date@0 as pickup_date], aggr=[count(Int64(1))] + FilterExec: pickup_date@0 = 2024-01-10 OR pickup_date@0 = 2024-01-11 + RepartitionExec: partitioning=RoundRobinBatch(16), input_partitions=1 + IcebergTableScan projection:[pickup_date] predicate:[(pickup_date = 2024-01-10) OR (pickup_date = 2024-01-11)] + "); + insta::assert_snapshot!(batches, @r" + +-------------+-------+ + | pickup_date | trips | + +-------------+-------+ + | 2024-01-10 | 25000 | + | 2024-01-11 | 25000 | + +-------------+-------+ + "); + + Ok(()) + } + + #[tokio::test] + async fn pushes_down_null_predicates() -> Result<()> { + let harness = IcebergTestHarness::new().await?; + let (plan, batches) = harness + .query("SELECT COUNT(*) AS trips FROM taxi WHERE passenger_count IS NULL") + .await?; + + insta::assert_snapshot!(plan, @r" + ProjectionExec: expr=[count(Int64(1))@0 as trips] + AggregateExec: mode=Final, gby=[], aggr=[count(Int64(1))] + CoalescePartitionsExec + AggregateExec: mode=Partial, gby=[], aggr=[count(Int64(1))] + FilterExec: passenger_count@0 IS NULL, projection=[] + RepartitionExec: partitioning=RoundRobinBatch(16), input_partitions=1 + IcebergTableScan projection:[passenger_count] predicate:[passenger_count IS NULL] + "); + insta::assert_snapshot!(batches, @r" + +-------+ + | trips | + +-------+ + | 8829 | + +-------+ + "); + + Ok(()) + } + + #[tokio::test] + async fn retains_unsupported_filter_after_safe_pushdown() -> Result<()> { + let harness = IcebergTestHarness::new().await?; + let (plan, batches) = harness + .query( + r"SELECT COUNT(*) AS trips + FROM taxi + WHERE pickup_date = DATE '2024-01-10' + AND trip_distance + 1.0 > 3.0", + ) + .await?; + + insta::assert_snapshot!(plan, @r" + ProjectionExec: expr=[count(Int64(1))@0 as trips] + AggregateExec: mode=Final, gby=[], aggr=[count(Int64(1))] + CoalescePartitionsExec + AggregateExec: mode=Partial, gby=[], aggr=[count(Int64(1))] + FilterExec: pickup_date@1 = 2024-01-10 AND trip_distance@0 + 1 > 3, projection=[] + RepartitionExec: partitioning=RoundRobinBatch(16), input_partitions=1 + IcebergTableScan projection:[trip_distance,pickup_date] predicate:[pickup_date = 2024-01-10] + "); + insta::assert_snapshot!(batches, @r" + +-------+ + | trips | + +-------+ + | 10516 | + +-------+ + "); + + Ok(()) + } + + #[tokio::test] + async fn executes_wholly_unsupported_filters_without_pushdown() -> Result<()> { + let harness = IcebergTestHarness::new().await?; + let (plan, batches) = harness + .query("SELECT COUNT(*) AS trips FROM taxi WHERE trip_distance + 1.0 > 20.0") + .await?; + + insta::assert_snapshot!(plan, @r" + ProjectionExec: expr=[count(Int64(1))@0 as trips] + AggregateExec: mode=Final, gby=[], aggr=[count(Int64(1))] + CoalescePartitionsExec + AggregateExec: mode=Partial, gby=[], aggr=[count(Int64(1))] + FilterExec: trip_distance@0 + 1 > 20, projection=[] + RepartitionExec: partitioning=RoundRobinBatch(16), input_partitions=1 + IcebergTableScan projection:[trip_distance] predicate:[] + "); + insta::assert_snapshot!(batches, @r" + +-------+ + | trips | + +-------+ + | 2757 | + +-------+ + "); + + Ok(()) + } +} diff --git a/crates/datafusion/tests/limit_pushdown.rs b/crates/datafusion/tests/limit_pushdown.rs new file mode 100644 index 0000000..ec19423 --- /dev/null +++ b/crates/datafusion/tests/limit_pushdown.rs @@ -0,0 +1,141 @@ +#[cfg(test)] +mod tests { + use datafusion::error::Result; + use datafusion_iceberg::test_utils::IcebergTestHarness; + + #[tokio::test] + async fn applies_sql_limit_to_the_iceberg_scan_output() -> Result<()> { + let harness = IcebergTestHarness::new().await?; + let (plan, batches) = harness + .query( + r"SELECT vendor_id, pickup_at, passenger_count, trip_distance, pickup_date + FROM taxi + WHERE pickup_date = DATE '2024-01-10' + ORDER BY pickup_at + LIMIT 3", + ) + .await?; + + insta::assert_snapshot!(plan, @r" + SortPreservingMergeExec: [pickup_at@1 ASC NULLS LAST], fetch=3 + SortExec: TopK(fetch=3), expr=[pickup_at@1 ASC NULLS LAST], preserve_partitioning=[true] + FilterExec: pickup_date@4 = 2024-01-10 + RepartitionExec: partitioning=RoundRobinBatch(16), input_partitions=1 + IcebergTableScan projection:[vendor_id,pickup_at,passenger_count,trip_distance,pickup_date] predicate:[pickup_date = 2024-01-10] + "); + insta::assert_snapshot!(batches, @r" + +-----------+---------------------+-----------------+---------------+-------------+ + | vendor_id | pickup_at | passenger_count | trip_distance | pickup_date | + +-----------+---------------------+-----------------+---------------+-------------+ + | 2 | 2024-01-10T00:00:09 | | 0.78 | 2024-01-10 | + | 2 | 2024-01-10T00:00:10 | 1 | 0.88 | 2024-01-10 | + | 1 | 2024-01-10T00:00:10 | 1 | 3.4 | 2024-01-10 | + +-----------+---------------------+-----------------+---------------+-------------+ + "); + + Ok(()) + } + + #[tokio::test] + async fn keeps_unordered_limits_above_the_iceberg_source() -> Result<()> { + let harness = IcebergTestHarness::new().await?; + let (plan, batches) = harness + .query("SELECT pickup_date FROM taxi WHERE pickup_date = DATE '2024-01-10' LIMIT 3") + .await?; + + insta::assert_snapshot!(plan, @r" + CoalescePartitionsExec: fetch=3 + FilterExec: pickup_date@0 = 2024-01-10, fetch=3 + RepartitionExec: partitioning=RoundRobinBatch(16), input_partitions=1 + IcebergTableScan projection:[pickup_date] predicate:[pickup_date = 2024-01-10] + "); + insta::assert_snapshot!(batches, @r" + +-------------+ + | pickup_date | + +-------------+ + | 2024-01-10 | + | 2024-01-10 | + | 2024-01-10 | + +-------------+ + "); + + Ok(()) + } + + #[tokio::test] + async fn passes_the_scan_limit_to_the_iceberg_source() -> Result<()> { + let harness = IcebergTestHarness::new().await?; + let (plan, batches) = harness + .query("SELECT pickup_date FROM taxi LIMIT 3") + .await?; + + insta::assert_snapshot!(plan, @r" + GlobalLimitExec: skip=0, fetch=3 + CooperativeExec + IcebergTableScan projection:[pickup_date] predicate:[] limit:[3] + "); + insta::assert_snapshot!(batches, @r" + +-------------+ + | pickup_date | + +-------------+ + | 2024-01-08 | + | 2024-01-08 | + | 2024-01-08 | + +-------------+ + "); + + Ok(()) + } + + #[tokio::test] + async fn avoids_scanning_for_limit_zero() -> Result<()> { + let harness = IcebergTestHarness::new().await?; + let (plan, batches) = harness + .query("SELECT vendor_id FROM taxi WHERE pickup_date = DATE '2024-01-10' LIMIT 0") + .await?; + + insta::assert_snapshot!(plan, @r" + EmptyExec + "); + insta::assert_snapshot!(batches, @r" + ++ + ++ + "); + + Ok(()) + } + + #[tokio::test] + async fn applies_offset_before_limit() -> Result<()> { + let harness = IcebergTestHarness::new().await?; + let (plan, batches) = harness + .query( + r"SELECT vendor_id, pickup_at + FROM taxi + WHERE pickup_date = DATE '2024-01-10' + ORDER BY pickup_at, vendor_id, pickup_location_id + LIMIT 2 OFFSET 3", + ) + .await?; + + insta::assert_snapshot!(plan, @r" + ProjectionExec: expr=[vendor_id@0 as vendor_id, pickup_at@1 as pickup_at] + GlobalLimitExec: skip=3, fetch=2 + SortPreservingMergeExec: [pickup_at@1 ASC NULLS LAST, vendor_id@0 ASC NULLS LAST, pickup_location_id@2 ASC NULLS LAST], fetch=5 + SortExec: TopK(fetch=5), expr=[pickup_at@1 ASC NULLS LAST, vendor_id@0 ASC NULLS LAST, pickup_location_id@2 ASC NULLS LAST], preserve_partitioning=[true] + FilterExec: pickup_date@3 = 2024-01-10, projection=[vendor_id@0, pickup_at@1, pickup_location_id@2] + RepartitionExec: partitioning=RoundRobinBatch(16), input_partitions=1 + IcebergTableScan projection:[vendor_id,pickup_at,pickup_location_id,pickup_date] predicate:[pickup_date = 2024-01-10] + "); + insta::assert_snapshot!(batches, @r" + +-----------+---------------------+ + | vendor_id | pickup_at | + +-----------+---------------------+ + | 2 | 2024-01-10T00:00:12 | + | 2 | 2024-01-10T00:00:14 | + +-----------+---------------------+ + "); + + Ok(()) + } +} diff --git a/crates/datafusion/tests/projection_pushdown.rs b/crates/datafusion/tests/projection_pushdown.rs new file mode 100644 index 0000000..4dd1eac --- /dev/null +++ b/crates/datafusion/tests/projection_pushdown.rs @@ -0,0 +1,140 @@ +#[cfg(test)] +mod tests { + use datafusion::error::Result; + use datafusion_iceberg::test_utils::IcebergTestHarness; + + #[tokio::test] + async fn projects_only_columns_required_by_the_query() -> Result<()> { + let harness = IcebergTestHarness::new().await?; + let (plan, batches) = harness + .query( + r"SELECT pickup_date, COUNT(*) AS trips + FROM taxi + WHERE pickup_date >= DATE '2024-01-10' + GROUP BY pickup_date + ORDER BY pickup_date", + ) + .await?; + + insta::assert_snapshot!(plan, @r" + SortPreservingMergeExec: [pickup_date@0 ASC NULLS LAST] + SortExec: expr=[pickup_date@0 ASC NULLS LAST], preserve_partitioning=[true] + ProjectionExec: expr=[pickup_date@0 as pickup_date, count(Int64(1))@1 as trips] + AggregateExec: mode=FinalPartitioned, gby=[pickup_date@0 as pickup_date], aggr=[count(Int64(1))] + RepartitionExec: partitioning=Hash([pickup_date@0], 16), input_partitions=16 + AggregateExec: mode=Partial, gby=[pickup_date@0 as pickup_date], aggr=[count(Int64(1))] + FilterExec: pickup_date@0 >= 2024-01-10 + RepartitionExec: partitioning=RoundRobinBatch(16), input_partitions=1 + IcebergTableScan projection:[pickup_date] predicate:[pickup_date >= 2024-01-10] + "); + insta::assert_snapshot!(batches, @r" + +-------------+-------+ + | pickup_date | trips | + +-------------+-------+ + | 2024-01-10 | 25000 | + | 2024-01-11 | 25000 | + | 2024-01-12 | 25000 | + | 2024-01-13 | 25000 | + | 2024-01-14 | 25000 | + +-------------+-------+ + "); + + Ok(()) + } + + #[tokio::test] + async fn includes_filter_columns_that_are_not_in_the_output() -> Result<()> { + let harness = IcebergTestHarness::new().await?; + let (plan, batches) = harness + .query( + r"SELECT vendor_id, pickup_location_id + FROM taxi + WHERE pickup_date = DATE '2024-01-10' + ORDER BY pickup_at, vendor_id + LIMIT 3", + ) + .await?; + + insta::assert_snapshot!(plan, @r" + ProjectionExec: expr=[vendor_id@0 as vendor_id, pickup_location_id@1 as pickup_location_id] + SortPreservingMergeExec: [pickup_at@2 ASC NULLS LAST, vendor_id@0 ASC NULLS LAST], fetch=3 + SortExec: TopK(fetch=3), expr=[pickup_at@2 ASC NULLS LAST, vendor_id@0 ASC NULLS LAST], preserve_partitioning=[true] + FilterExec: pickup_date@3 = 2024-01-10, projection=[vendor_id@0, pickup_location_id@2, pickup_at@1] + RepartitionExec: partitioning=RoundRobinBatch(16), input_partitions=1 + IcebergTableScan projection:[vendor_id,pickup_at,pickup_location_id,pickup_date] predicate:[pickup_date = 2024-01-10] + "); + insta::assert_snapshot!(batches, @r" + +-----------+--------------------+ + | vendor_id | pickup_location_id | + +-----------+--------------------+ + | 2 | 75 | + | 1 | 161 | + | 2 | 162 | + +-----------+--------------------+ + "); + + Ok(()) + } + + #[tokio::test] + async fn projects_source_columns_for_computed_expressions() -> Result<()> { + let harness = IcebergTestHarness::new().await?; + let (plan, batches) = harness + .query( + r"SELECT MAX(trip_distance * fare_amount) AS max_weighted_fare + FROM taxi + WHERE pickup_date = DATE '2024-01-10'", + ) + .await?; + + insta::assert_snapshot!(plan, @r" + ProjectionExec: expr=[max(taxi.trip_distance * taxi.fare_amount)@0 as max_weighted_fare] + AggregateExec: mode=Final, gby=[], aggr=[max(taxi.trip_distance * taxi.fare_amount)] + CoalescePartitionsExec + AggregateExec: mode=Partial, gby=[], aggr=[max(taxi.trip_distance * taxi.fare_amount)] + FilterExec: pickup_date@2 = 2024-01-10, projection=[trip_distance@0, fare_amount@1] + RepartitionExec: partitioning=RoundRobinBatch(16), input_partitions=1 + IcebergTableScan projection:[trip_distance,fare_amount,pickup_date] predicate:[pickup_date = 2024-01-10] + "); + insta::assert_snapshot!(batches, @r" + +-------------------+ + | max_weighted_fare | + +-------------------+ + | 207508.9584 | + +-------------------+ + "); + + Ok(()) + } + + #[tokio::test] + async fn reads_all_columns_when_selecting_star() -> Result<()> { + let harness = IcebergTestHarness::new().await?; + let (plan, batches) = harness + .query( + r"SELECT * + FROM taxi + WHERE pickup_date = DATE '2024-01-10' + ORDER BY pickup_at, vendor_id, pickup_location_id + LIMIT 1", + ) + .await?; + + insta::assert_snapshot!(plan, @r" + SortPreservingMergeExec: [pickup_at@1 ASC NULLS LAST, vendor_id@0 ASC NULLS LAST, pickup_location_id@5 ASC NULLS LAST], fetch=1 + SortExec: TopK(fetch=1), expr=[pickup_at@1 ASC NULLS LAST, vendor_id@0 ASC NULLS LAST, pickup_location_id@5 ASC NULLS LAST], preserve_partitioning=[true] + FilterExec: pickup_date@12 = 2024-01-10 + RepartitionExec: partitioning=RoundRobinBatch(16), input_partitions=1 + IcebergTableScan projection:[vendor_id,pickup_at,dropoff_at,passenger_count,trip_distance,pickup_location_id,dropoff_location_id,payment_type,fare_amount,tip_amount,tolls_amount,total_amount,pickup_date] predicate:[pickup_date = 2024-01-10] + "); + insta::assert_snapshot!(batches, @r" + +-----------+---------------------+---------------------+-----------------+---------------+--------------------+---------------------+--------------+-------------+------------+--------------+--------------+-------------+ + | vendor_id | pickup_at | dropoff_at | passenger_count | trip_distance | pickup_location_id | dropoff_location_id | payment_type | fare_amount | tip_amount | tolls_amount | total_amount | pickup_date | + +-----------+---------------------+---------------------+-----------------+---------------+--------------------+---------------------+--------------+-------------+------------+--------------+--------------+-------------+ + | 2 | 2024-01-10T00:00:09 | 2024-01-10T00:03:30 | | 0.78 | 75 | 236 | 0 | 1.74 | 3.15 | 0.0 | 8.89 | 2024-01-10 | + +-----------+---------------------+---------------------+-----------------+---------------+--------------------+---------------------+--------------+-------------+------------+--------------+--------------+-------------+ + "); + + Ok(()) + } +} From 3e73d8f8c8fffeaed12514081033a4cac05cd936 Mon Sep 17 00:00:00 2001 From: Gabriel Date: Wed, 9 Sep 2026 14:00:24 +0200 Subject: [PATCH 2/5] test: update DataFusion 55 plan snapshots --- crates/datafusion/tests/filter_pushdown.rs | 4 ++-- crates/datafusion/tests/projection_pushdown.rs | 4 ++-- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/crates/datafusion/tests/filter_pushdown.rs b/crates/datafusion/tests/filter_pushdown.rs index d82f107..af13025 100644 --- a/crates/datafusion/tests/filter_pushdown.rs +++ b/crates/datafusion/tests/filter_pushdown.rs @@ -52,8 +52,8 @@ mod tests { insta::assert_snapshot!(plan, @r" SortPreservingMergeExec: [pickup_date@0 ASC NULLS LAST] - SortExec: expr=[pickup_date@0 ASC NULLS LAST], preserve_partitioning=[true] - ProjectionExec: expr=[pickup_date@0 as pickup_date, count(Int64(1))@1 as trips] + ProjectionExec: expr=[pickup_date@0 as pickup_date, count(Int64(1))@1 as trips] + SortExec: expr=[pickup_date@0 ASC NULLS LAST], preserve_partitioning=[true] AggregateExec: mode=FinalPartitioned, gby=[pickup_date@0 as pickup_date], aggr=[count(Int64(1))] RepartitionExec: partitioning=Hash([pickup_date@0], 16), input_partitions=16 AggregateExec: mode=Partial, gby=[pickup_date@0 as pickup_date], aggr=[count(Int64(1))] diff --git a/crates/datafusion/tests/projection_pushdown.rs b/crates/datafusion/tests/projection_pushdown.rs index 4dd1eac..0cd1aed 100644 --- a/crates/datafusion/tests/projection_pushdown.rs +++ b/crates/datafusion/tests/projection_pushdown.rs @@ -18,8 +18,8 @@ mod tests { insta::assert_snapshot!(plan, @r" SortPreservingMergeExec: [pickup_date@0 ASC NULLS LAST] - SortExec: expr=[pickup_date@0 ASC NULLS LAST], preserve_partitioning=[true] - ProjectionExec: expr=[pickup_date@0 as pickup_date, count(Int64(1))@1 as trips] + ProjectionExec: expr=[pickup_date@0 as pickup_date, count(Int64(1))@1 as trips] + SortExec: expr=[pickup_date@0 ASC NULLS LAST], preserve_partitioning=[true] AggregateExec: mode=FinalPartitioned, gby=[pickup_date@0 as pickup_date], aggr=[count(Int64(1))] RepartitionExec: partitioning=Hash([pickup_date@0], 16), input_partitions=16 AggregateExec: mode=Partial, gby=[pickup_date@0 as pickup_date], aggr=[count(Int64(1))] From 9c93f1506d8bb1b2b86448ae6fe2f0d541e1b501 Mon Sep 17 00:00:00 2001 From: Gabriel Date: Fri, 25 Sep 2026 13:24:12 +0200 Subject: [PATCH 3/5] ci: fetch LFS test fixtures --- .github/workflows/ci.yml | 1 + 1 file changed, 1 insertion(+) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 89c0545..f9184bc 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -76,6 +76,7 @@ jobs: steps: - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1 with: + lfs: true persist-credentials: false - name: Setup Rust uses: ./.github/actions/setup-rust From da4618732ba9bcf4750fdf20ecee6023cc912d51 Mon Sep 17 00:00:00 2001 From: Gabriel Date: Fri, 25 Sep 2026 13:25:14 +0200 Subject: [PATCH 4/5] docs: explain LFS fixture setup --- crates/datafusion/testdata/iceberg/taxi/README.md | 13 +++++++++++++ 1 file changed, 13 insertions(+) diff --git a/crates/datafusion/testdata/iceberg/taxi/README.md b/crates/datafusion/testdata/iceberg/taxi/README.md index b4fed84..af82215 100644 --- a/crates/datafusion/testdata/iceberg/taxi/README.md +++ b/crates/datafusion/testdata/iceberg/taxi/README.md @@ -9,6 +9,19 @@ The table uses the stable URI prefix `s3://iceberg-test/warehouse/taxi/` in its metadata. Tests map that prefix to this checked-in directory, so the metadata stays valid regardless of where the repository is cloned. +## Git LFS setup + +The Parquet data files are tracked with Git LFS. Install Git LFS and run the +following before cloning the repository so the fixture files are fetched rather +than their LFS pointer files: + +```shell +git lfs install +``` + +For an existing clone made before Git LFS was configured, run `git lfs pull` to +fetch the fixture files. + ## Selection - Source: `yellow_tripdata_2024-01.parquet`, published by the NYC TLC. From 8215aa05a0037f2fba6b05b0c8b3b857134c687b Mon Sep 17 00:00:00 2001 From: Gabriel Date: Fri, 25 Sep 2026 13:51:21 +0200 Subject: [PATCH 5/5] test: stabilize fixture plan snapshots --- crates/datafusion/src/test_utils/harness.rs | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/crates/datafusion/src/test_utils/harness.rs b/crates/datafusion/src/test_utils/harness.rs index 90e93c0..39ea8d2 100644 --- a/crates/datafusion/src/test_utils/harness.rs +++ b/crates/datafusion/src/test_utils/harness.rs @@ -8,7 +8,7 @@ use datafusion::dataframe::DataFrame; use datafusion::error::Result; use datafusion::execution::SessionStateBuilder; use datafusion::physical_plan::displayable; -use datafusion::prelude::SessionContext; +use datafusion::prelude::{SessionConfig, SessionContext}; use futures::StreamExt; use futures::stream::BoxStream; use iceberg::io::{ @@ -31,6 +31,8 @@ impl IcebergTestHarness { pub async fn new() -> Result { let state = SessionStateBuilder::new() .with_default_features() + // Plan snapshots must not vary with the runner's CPU count. + .with_config(SessionConfig::new().with_target_partitions(16)) .with_table_factory( "ICEBERG".to_string(), Arc::new(IcebergTableProviderFactory::new_with_storage_factory(