use super::*;
use crate::formats::schema_union::partition_columns_of_key as partition_columns_from_prefix;
use crate::formats::schema_union::{FOOTERS_AT_ONCE, lenient_scan, with_partition_columns};
use polars::prelude::{NamedFrom, Series};
async fn list_dataset_files(
store: &Arc<dyn ObjectStore>,
prefix: &str,
pattern: Option<&globset::GlobMatcher>,
) -> Result<(Vec<DatasetFile>, crate::formats::schema_union::SkippedFiles)> {
let progress = crate::formats::schema_union::FooterProgress::default();
let listing = progress.listing();
list_dataset_files_reporting(
store,
prefix,
pattern,
ListShards::ONE,
listing.counter(),
listing.cancel_flag(),
)
.await
}
fn dataset_schema_from_footers(
files: &[DatasetFile],
read: &[usize],
footers: &[Option<FileFooter>],
) -> Result<(crate::formats::schema_union::DatasetSchema, Vec<String>)> {
let (first, newest) = match files {
[] => {
return Err(color_eyre::eyre::eyre!(
"No parquet file found in cloud prefix"
));
}
[only] => (only, only),
[first, .., last] => (first, last),
};
let mut union = crate::formats::schema_union::union_sampled(files.len(), read, footers);
if union.schema.is_empty() {
return Err(color_eyre::eyre::eyre!(
"No readable parquet footer in cloud prefix"
));
}
let (partition_columns, values) =
crate::formats::schema_union::partitions_of_listing(&first.key, &newest.key);
union.schema = Arc::new(with_partition_columns(
&union.schema,
&partition_columns,
&values,
));
Ok((union, partition_columns))
}
async fn footers_of_files(
store: &Arc<dyn ObjectStore>,
files: &[DatasetFile],
read: &[usize],
meter: &Arc<crate::loading::measurements::Meter>,
) -> Vec<Option<FileFooter>> {
footers_of_files_reporting(
store,
files,
read,
&crate::formats::schema_union::FooterProgress::default(),
meter,
)
.await
}
#[test]
fn partition_columns_from_prefix_basic() {
let cols = partition_columns_from_prefix("dataset/year=2024/month=01");
assert_eq!(cols, ["year", "month"]);
}
#[test]
fn partition_columns_from_prefix_with_trailing_slash() {
let cols = partition_columns_from_prefix("path/year=2024/month=01/day=15/");
assert_eq!(cols, ["year", "month", "day"]);
}
#[test]
fn partition_columns_from_prefix_dedup() {
let cols = partition_columns_from_prefix("a/x=1/x=2");
assert_eq!(cols, ["x"]);
}
#[test]
fn partition_columns_from_prefix_empty() {
let cols = partition_columns_from_prefix("");
assert!(cols.is_empty());
}
async fn schema_of(
store: &Arc<dyn ObjectStore>,
files: &[DatasetFile],
) -> (crate::formats::schema_union::DatasetSchema, Vec<String>) {
let read: Vec<usize> = (0..files.len()).collect();
let footers = footers_of_files(
store,
files,
&read,
&Arc::new(crate::loading::measurements::Meter::default()),
)
.await;
dataset_schema_from_footers(files, &read, &footers).unwrap()
}
fn evolving_dataset() -> Vec<(String, Vec<u8>)> {
use polars::prelude::{IntoSeries, ParquetWriter, StructChunked, df};
let write = |mut df: polars::prelude::DataFrame| {
let mut bytes = Vec::new();
ParquetWriter::new(&mut bytes).finish(&mut df).unwrap();
bytes
};
let old_input = StructChunked::from_series(
"input".into(),
2,
[Series::new("value".into(), &[1.0f64, 2.0])].iter(),
)
.unwrap()
.into_series();
let old = df!("id" => &[1i64, 2], "input" => old_input).unwrap();
let new_input = StructChunked::from_series(
"input".into(),
5,
[
Series::new("value".into(), &[3.0f64, 4.0, 5.0, 6.0, 7.0]),
Series::new("address".into(), &["a", "b", "c", "d", "e"]),
]
.iter(),
)
.unwrap()
.into_series();
let new =
df!("id" => &[3i64, 4, 5, 6, 7], "fee" => &[10i64, 20, 30, 40, 50], "input" => new_input)
.unwrap();
vec![
(
"data/date=2009-01-03/part-0.parquet".to_string(),
write(old),
),
(
"data/date=2026-09-17/part-0.parquet".to_string(),
write(new),
),
("data/date=2026-09-17/_SUCCESS".to_string(), b"x".to_vec()),
(
"data/date=2026-09-17/.part-0.parquet.crc".to_string(),
b"crc".to_vec(),
),
]
}
#[test]
fn split_points_follow_the_keys_down() {
let prefix = "parquet/by_station/";
let first = "parquet/by_station/STATION=ACW00011604/ELEMENT=PGTM/a.parquet";
let last = "parquet/by_station/STATION=AEM00041217/ELEMENT=TMAX/b.parquet";
let points = split_points(first, last, None, prefix.len(), 3);
assert_eq!(points.len(), 3, "{points:?}");
assert!(points.windows(2).all(|w| w[0] < w[1]), "sorted, no repeats");
assert!(points.iter().all(|p| p.as_str() > last));
for point in &points {
let id = point
.strip_prefix("parquet/by_station/STATION=")
.unwrap_or_else(|| panic!("{point}"));
assert!(id.starts_with(|c: char| c.is_ascii_uppercase()), "{point}");
}
let first = "parquet/by_station/STATION=US009052008/ELEMENT=PRCP/a.parquet";
let last = "parquet/by_station/STATION=US1AKAB0001/ELEMENT=PRCP/a.parquet";
let through = "parquet/by_station/STATION=US2";
let points = split_points(first, last, Some(through), prefix.len(), 4);
assert_eq!(points.len(), 4, "{points:?}");
assert!(points.windows(2).all(|w| w[0] < w[1]));
for point in &points {
assert!(point.as_str() > last && point.as_str() < through, "{point}");
assert!(
point.starts_with("parquet/by_station/STATION=US1"),
"{point}"
);
}
assert!(split_points(first, last, Some(through), prefix.len(), 0).is_empty());
assert!(split_points(first, last, Some(&format!("{last}0")), prefix.len(), 4).is_empty());
}
fn skewed_keys() -> Vec<String> {
let mut keys = Vec::new();
let elements = ["PRCP", "SNOW", "TMAX"];
let mut station = |code: String| {
for element in elements {
keys.push(format!(
"p/by_station/STATION={code}/ELEMENT={element}/x_0.snappy.parquet"
));
}
};
for i in 0..60 {
station(format!("AC{i:09}"));
}
for country in ["BR", "CA", "GM", "SF", "UK"] {
for i in 0..40 {
station(format!("{country}{i:09}"));
}
}
for kind in ["1AK", "1CA", "1TX", "C00", "W00"] {
for i in 0..300 {
station(format!("US{kind}{i:06}"));
}
}
for odd in [
"p/by_station/STATION=B",
"p/by_station/STATION=U",
"p/by_station/STATION=US1B",
"p/by_station/STATION=Z~tilde/a.parquet",
"p/by_station/STATION=ü/a.parquet",
"p/by_station/_SUCCESS",
"p/by_station/zz/a.parquet",
] {
keys.push(odd.to_string());
}
keys
}
#[test]
fn a_listing_in_ranges_finds_what_one_listing_does() {
use object_store::PutPayload;
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
let keys = skewed_keys();
rt.block_on(async {
for key in keys.iter().map(String::as_str).chain([
"p/by_stationx/a.parquet",
"p/a.parquet",
"q/b.parquet",
]) {
store
.put(&OsPath::from(key), PutPayload::from(b"x".to_vec()))
.await
.unwrap();
}
let prefix = OsPath::from("p/by_station");
let (listed, cancelled): (
Arc<std::sync::atomic::AtomicUsize>,
Arc<std::sync::atomic::AtomicBool>,
) = (Arc::default(), Arc::default());
let keys_of = |objects: Vec<object_store::ObjectMeta>| {
let mut keys: Vec<String> = objects
.into_iter()
.map(|o| o.location.as_ref().to_string())
.collect();
keys.sort();
keys
};
let one = keys_of(
list_objects(
&store,
Some(&prefix),
ListShards::ONE,
listed.clone(),
cancelled.clone(),
)
.await
.unwrap(),
);
let mut expected: Vec<String> = keys
.iter()
.map(|k| OsPath::from(k.as_str()).to_string())
.collect();
expected.sort();
assert_eq!(one, expected);
for (at_once, most, split_after, split_into) in [
(2, 8, 10, 1),
(8, 64, 25, 4),
(64, 512, 7, 64),
(4, 1000, 1, 2),
(64, 1024, 30, 4),
] {
let plan = ListShards {
at_once,
most,
split_after,
split_into,
};
let ranges = keys_of(
list_objects(
&store,
Some(&prefix),
plan,
listed.clone(),
cancelled.clone(),
)
.await
.unwrap(),
);
assert_eq!(ranges.len(), one.len(), "{plan:?}: a key twice or missing");
assert_eq!(ranges, one, "{plan:?}");
}
let (whole, skipped) = list_dataset_files(&store, "p/by_station", None)
.await
.unwrap();
let (shared, shared_skipped) = list_dataset_files_reporting(
&store,
"p/by_station",
None,
ListShards {
at_once: 8,
most: 64,
split_after: 20,
split_into: 4,
},
listed,
cancelled,
)
.await
.unwrap();
assert_eq!(whole, shared);
assert_eq!(skipped, shared_skipped);
});
}
#[test]
fn only_stores_that_list_from_an_offset_list_in_ranges() {
let parallel = |url: &str| ListShards::for_url(url).at_once > 1;
assert!(parallel("s3://noaa-ghcn-pds/parquet/by_station/"));
assert!(parallel("gs://bucket/data/"));
assert!(!parallel("s3://my-bucket--usw2-az1--x-s3/data/"));
assert!(!parallel("az://container/data/"));
assert!(!parallel("https://example.com/data/"));
}
#[test]
fn a_listing_counts_what_it_finds() {
use object_store::PutPayload;
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
rt.block_on(async {
for key in ["data/a.parquet", "data/b.parquet", "data/_SUCCESS"] {
store
.put(&OsPath::from(key), PutPayload::from(b"PAR1".to_vec()))
.await
.unwrap();
}
let progress = crate::formats::schema_union::FooterProgress::default();
{
let listing = progress.listing();
let (files, _skipped) = list_dataset_files_reporting(
&store,
"data/",
None,
ListShards::ONE,
listing.counter(),
listing.cancel_flag(),
)
.await
.unwrap();
assert_eq!(files.len(), 2);
assert_eq!(progress.listed(), Some(3), "every object, data or not");
}
assert_eq!(progress.listed(), None, "a finished listing shows no count");
let listing = progress.listing();
assert_eq!(
progress.listed(),
Some(0),
"a new listing starts from nothing"
);
progress.cancel();
assert!(
list_dataset_files_reporting(
&store,
"data/",
None,
ListShards::ONE,
listing.counter(),
listing.cancel_flag(),
)
.await
.is_err(),
"an abandoned load stops listing"
);
});
}
#[test]
fn a_cancelled_pass_reads_no_footers() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let mut frame = df!("n" => (0..10i64).collect::<Vec<i64>>()).unwrap();
let mut bytes = Vec::new();
ParquetWriter::new(&mut bytes).finish(&mut frame).unwrap();
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
rt.block_on(async {
for key in ["data/a.parquet", "data/b.parquet"] {
store
.put(&OsPath::from(key), PutPayload::from(bytes.clone()))
.await
.unwrap();
}
let (files, _skipped) = list_dataset_files(&store, "data/", None).await.unwrap();
let read: Vec<usize> = (0..files.len()).collect();
let progress = crate::formats::schema_union::FooterProgress::default();
progress.cancel();
let meter = Arc::new(crate::loading::measurements::Meter::default());
let footers = footers_of_files_reporting(&store, &files, &read, &progress, &meter).await;
assert!(footers.iter().all(|f| f.is_none()), "nothing was read");
let requests = meter
.footers()
.and_then(|cost| cost.over_the_wire)
.map(|wire| wire.requests)
.unwrap_or(0);
assert_eq!(requests, 0, "nothing was requested");
});
}
#[test]
fn a_remote_dataset_carries_its_row_group_sizes_into_the_schema() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let write = |groups: i64| -> Vec<u8> {
let mut frame = df!("n" => (0..groups * 1_000).collect::<Vec<i64>>()).unwrap();
let mut bytes = Vec::new();
ParquetWriter::new(&mut bytes)
.with_row_group_size(Some(1_000))
.finish(&mut frame)
.unwrap();
bytes
};
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
rt.block_on(async {
for (key, bytes) in [
("data/date=2024-01-01/a.parquet", write(3)),
("data/date=2024-01-02/b.parquet", write(2)),
] {
store
.put(&OsPath::from(key), PutPayload::from(bytes))
.await
.unwrap();
}
let (files, _skipped) = list_dataset_files(&store, "data/", None).await.unwrap();
let read: Vec<usize> = (0..files.len()).collect();
let footers = footers_of_files(
&store,
&files,
&read,
&Arc::new(crate::loading::measurements::Meter::default()),
)
.await;
let mut sizes: Vec<usize> = footers
.iter()
.flatten()
.flat_map(|footer| footer.row_group_bytes.iter().copied())
.collect();
assert_eq!(sizes.len(), 5, "three row groups and two: {sizes:?}");
assert!(sizes.iter().all(|size| *size > 0), "{sizes:?}");
let (dataset, _) = dataset_schema_from_footers(&files, &read, &footers).unwrap();
sizes.sort_unstable();
assert_eq!(
dataset.median_row_group_bytes,
Some(sizes[2]),
"the middle of the five reaches the schema: {sizes:?}"
);
let rows: Vec<String> = (0..20_000)
.map(|i| format!("{i:0>6}{}", "abcdefghij".repeat(19)))
.collect();
let mut wide = df!("s" => rows).unwrap();
let mut bytes = Vec::new();
ParquetWriter::new(&mut bytes)
.with_row_group_size(Some(20_000))
.finish(&mut wide)
.unwrap();
store
.put(
&OsPath::from("wide/date=2024-01-01/w.parquet"),
PutPayload::from(bytes),
)
.await
.unwrap();
let (wide_files, _skipped) = list_dataset_files(&store, "wide/", None).await.unwrap();
assert_eq!(wide_files.len(), 1, "only the wide file: {wide_files:?}");
let wide_footers = footers_of_files(
&store,
&wide_files,
&[0],
&Arc::new(crate::loading::measurements::Meter::default()),
)
.await;
let size = wide_footers[0].as_ref().unwrap().row_group_bytes[0];
assert!(
size < 1_000_000,
"the compressed size, not the decoded one: {size} bytes"
);
});
}
#[test]
fn a_sampled_remote_read_takes_each_size_from_the_file_it_read() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let write = |rows: i64| -> Vec<u8> {
let mut frame = df!("n" => (0..rows).collect::<Vec<i64>>()).unwrap();
let mut bytes = Vec::new();
ParquetWriter::new(&mut bytes).finish(&mut frame).unwrap();
bytes
};
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
rt.block_on(async {
for (key, rows) in [
("s/date=2024-01-01/a.parquet", 1),
("s/date=2024-01-02/b.parquet", 200),
("s/date=2024-01-03/c.parquet", 40_000),
] {
store
.put(&OsPath::from(key), PutPayload::from(write(rows)))
.await
.unwrap();
}
let (files, _skipped) = list_dataset_files(&store, "s/", None).await.unwrap();
assert_eq!(files.len(), 3);
let read = [2usize];
let footers = footers_of_files(
&store,
&files,
&read,
&Arc::new(crate::loading::measurements::Meter::default()),
)
.await;
let (dataset, _) = dataset_schema_from_footers(&files, &read, &footers).unwrap();
assert_eq!(
dataset.median_file_bytes,
Some(files[2].size as usize),
"the file read, not the first in the list: {:?}",
files.iter().map(|f| f.size).collect::<Vec<_>>()
);
});
}
#[test]
fn a_cloud_open_counts_its_footers_against_the_counter_it_is_given() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let write = |rows: i64| -> Vec<u8> {
let mut frame = df!("n" => (0..rows).collect::<Vec<i64>>()).unwrap();
let mut bytes = Vec::new();
ParquetWriter::new(&mut bytes).finish(&mut frame).unwrap();
bytes
};
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
rt.block_on(async {
for (key, rows) in [
("data/date=2024-01-01/a.parquet", 1),
("data/date=2024-01-02/b.parquet", 2),
("data/date=2024-01-03/c.parquet", 3),
] {
store
.put(&OsPath::from(key), PutPayload::from(write(rows)))
.await
.unwrap();
}
});
let progress = Arc::new(crate::formats::schema_union::FooterProgress::default());
let _ = crate::App::schema_state_from_cloud_hive_with(
"memory://data/".to_string(),
"data/".to_string(),
store,
polars::prelude::cloud::CloudOptions::default(),
&crate::OpenOptions::default(),
rt.handle(),
&crate::loading::measurements::OpenReport {
progress: progress.clone(),
meter: Arc::new(crate::loading::measurements::Meter::default()),
remembered: None,
writes: Default::default(),
},
);
let pass = progress.last_pass();
assert_eq!(pass.begun, 1, "the open ran its footer pass against it");
assert_eq!(pass.read, 3, "counting each of the three objects off");
assert_eq!(
progress.reading(),
None,
"with nothing left to say once they landed"
);
}
#[test]
fn a_cloud_count_that_read_nothing_leaves_the_measurement_for_the_one_that_does() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
let key = OsPath::from("data/date=2024-01-01/half-written.parquet");
rt.block_on(async {
store
.put(&key, PutPayload::from(b"not parquet yet".to_vec()))
.await
.unwrap();
});
let files = vec![DatasetFile {
key: "data/date=2024-01-01/half-written.parquet".to_string(),
size: 15,
stamp: 0,
etag: None,
}];
let meter = Arc::new(crate::loading::measurements::Meter::default());
meter.listed(std::time::Duration::from_millis(1), Some(1), false);
let source = crate::formats::dataset_files::StoreFiles::new(
"memory://data/",
"data/".to_string(),
None,
store.clone(),
polars::prelude::cloud::CloudOptions::default(),
rt.handle(),
);
crate::formats::dataset_files::footers_for_count(
&source,
&Arc::new(files),
&[0],
&meter,
&Default::default(),
);
assert_eq!(
meter.footers(),
None,
"nothing under there parsed, so nothing was measured and the one \
measurement counting gets is still to be had"
);
let mut frame = df!("n" => &[1i64, 2, 3]).unwrap();
let mut body = Vec::new();
ParquetWriter::new(&mut body).finish(&mut frame).unwrap();
let size = body.len() as u64;
rt.block_on(async {
store.put(&key, PutPayload::from(body)).await.unwrap();
});
let files = vec![DatasetFile {
key: "data/date=2024-01-01/half-written.parquet".to_string(),
size,
stamp: 0,
etag: None,
}];
let footers = crate::formats::dataset_files::footers_for_count(
&source,
&Arc::new(files),
&[0],
&meter,
&Default::default(),
)
.unwrap();
assert_eq!(
footers
.iter()
.flatten()
.flat_map(|f| f.row_group_rows.iter())
.sum::<usize>(),
3,
"and it counts the three rows"
);
let footers = meter.footers().expect("the count that worked was measured");
assert_eq!(footers.files, Some(1), "over the one object");
assert!(
footers
.over_the_wire
.is_some_and(|w| w.bytes.is_some_and(|b| b > 15)),
"reporting the real read, not the fifteen bytes of the half-written one; \
got {:?}",
footers.over_the_wire
);
}
#[test]
fn a_footer_pass_survives_the_cache_and_comes_back_the_same() {
use polars::prelude::DataType;
let schema_of = |cols: &[(&str, DataType)]| {
let mut schema = Schema::with_capacity(cols.len());
for (name, dtype) in cols {
schema.with_column((*name).into(), dtype.clone());
}
Arc::new(schema)
};
let original = vec![
Some(FileFooter {
schema: schema_of(&[("id", DataType::Int64), ("note", DataType::String)]),
row_group_rows: vec![100, 50],
row_group_bytes: vec![4_096, 2_048],
file_bytes: 10,
column_bytes: vec![("id".into(), 1_200), ("note".into(), 900)],
}),
Some(FileFooter {
schema: schema_of(&[("id", DataType::Int64), ("note", DataType::String)]),
row_group_rows: vec![7],
row_group_bytes: vec![512],
file_bytes: 20,
column_bytes: Vec::new(),
}),
Some(FileFooter {
schema: schema_of(&[("id", DataType::Int64), ("extra", DataType::Boolean)]),
row_group_rows: vec![3],
row_group_bytes: vec![128],
file_bytes: 30,
column_bytes: Vec::new(),
}),
None,
];
let (cached, schemas) = crate::formats::schema_union::footers_to_cache(&original);
let sizes = [10, 20, 30, 0];
assert_eq!(
schemas.len(),
2,
"two distinct shapes among four files, not four copies of them"
);
assert_eq!(cached[3].schema, None, "and the unreadable one says so");
let back = crate::formats::schema_union::footers_from_cache(&cached, &schemas, &sizes)
.expect("the table is consistent");
assert_eq!(back.len(), original.len());
for (before, after) in original.iter().zip(&back) {
match (before, after) {
(None, None) => {}
(Some(a), Some(b)) => {
assert_eq!(a.schema, b.schema, "same columns, same types");
assert_eq!(a.row_group_rows, b.row_group_rows);
assert_eq!(a.row_group_bytes, b.row_group_bytes);
assert_eq!(a.file_bytes, b.file_bytes, "the size, from the listing");
assert_eq!(a.column_bytes, b.column_bytes);
}
_ => panic!("a footer changed whether it could be read"),
}
}
let broken = vec![crate::cache::CachedFooter {
schema: Some(9),
row_group_rows: vec![1],
row_group_bytes: vec![1],
column_bytes: Vec::new(),
}];
assert!(
crate::formats::schema_union::footers_from_cache(&broken, &schemas, &[1]).is_none(),
"an entry that points at a schema it does not have is refused whole"
);
}
#[test]
fn a_dataset_read_behind_the_open_is_remembered_by_the_pass_that_read_it() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let body = || {
let mut frame = df!("n" => &[1i64]).unwrap();
let mut out = Vec::new();
ParquetWriter::new(&mut out).finish(&mut frame).unwrap();
out
};
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
let count = FOOTERS_AT_ONCE + 1;
rt.block_on(async {
for i in 0..count {
store
.put(
&OsPath::from(format!("data/f{i:04}.parquet")),
PutPayload::from(body()),
)
.await
.unwrap();
}
});
let dir = tempfile::tempdir().unwrap();
let cache = crate::cache::CacheManager::with_dir(dir.path().to_path_buf());
let report = || crate::loading::measurements::OpenReport {
progress: Arc::new(crate::formats::schema_union::FooterProgress::default()),
meter: Arc::new(crate::loading::measurements::Meter::default()),
remembered: Some(cache.clone()),
writes: Default::default(),
};
let open = || {
let r = report();
let state = crate::App::schema_state_from_cloud_hive_with(
"memory://data/".to_string(),
"data/".to_string(),
store.clone(),
polars::prelude::cloud::CloudOptions::default(),
&crate::OpenOptions::default(),
rt.handle(),
&r,
);
(state, r.meter.clone())
};
let (state, meter) = open();
assert_eq!(
meter.footers().and_then(|c| c.files),
Some(2),
"the open itself reads the two ends, which is what staging is"
);
assert!(
cache.dataset_shapes_kept() == 0,
"and keeps nothing: a two-footer view of {count} files is not this dataset, \
and kept as one it would open next time with two files' worth of schema"
);
let pending = state
.and_then(|(_, facts)| facts.footers_pending)
.expect("a staged open leaves a pass behind it");
let _ = pending(&Arc::new(
crate::formats::schema_union::FooterProgress::default(),
));
let (_, second) = open();
assert_eq!(
second.footers(),
None,
"so the next open reads no footers at all, for a dataset of {count} files"
);
}
#[test]
fn a_footer_that_would_not_read_is_not_remembered_as_unreadable() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let good = || {
let mut frame = df!("n" => &[1i64]).unwrap();
let mut out = Vec::new();
ParquetWriter::new(&mut out).finish(&mut frame).unwrap();
out
};
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
rt.block_on(async {
store
.put(&OsPath::from("data/a.parquet"), PutPayload::from(good()))
.await
.unwrap();
store
.put(
&OsPath::from("data/b.parquet"),
PutPayload::from(b"not parquet yet".to_vec()),
)
.await
.unwrap();
});
let dir = tempfile::tempdir().unwrap();
let cache = crate::cache::CacheManager::with_dir(dir.path().to_path_buf());
let open = || {
let meter = Arc::new(crate::loading::measurements::Meter::default());
let _ = crate::App::schema_state_from_cloud_hive_with(
"memory://data/".to_string(),
"data/".to_string(),
store.clone(),
polars::prelude::cloud::CloudOptions::default(),
&crate::OpenOptions::default(),
rt.handle(),
&crate::loading::measurements::OpenReport {
progress: Arc::new(crate::formats::schema_union::FooterProgress::default()),
meter: meter.clone(),
remembered: Some(cache.clone()),
writes: Default::default(),
},
);
meter
};
let first = open();
assert_eq!(
first.footers().and_then(|c| c.files),
Some(2),
"both were tried"
);
assert!(
cache.dataset_shapes_kept() == 0,
"and nothing was kept, because one of them did not come back"
);
let second = open();
assert_eq!(
second.footers().and_then(|c| c.files),
Some(2),
"so the next open tries again rather than taking the failure as settled"
);
}
#[test]
fn a_dataset_opened_again_is_not_read_again() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let body = |rows: i64| {
let mut frame = df!("n" => (0..rows).collect::<Vec<i64>>()).unwrap();
let mut out = Vec::new();
ParquetWriter::new(&mut out).finish(&mut frame).unwrap();
out
};
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
rt.block_on(async {
for (key, rows) in [("data/a.parquet", 3i64), ("data/b.parquet", 4)] {
store
.put(&OsPath::from(key), PutPayload::from(body(rows)))
.await
.unwrap();
}
});
let dir = tempfile::tempdir().unwrap();
let cache = crate::cache::CacheManager::with_dir(dir.path().to_path_buf());
let open = |store: Arc<dyn ObjectStore>| {
let meter = Arc::new(crate::loading::measurements::Meter::default());
let _ = crate::App::schema_state_from_cloud_hive_with(
"memory://data/".to_string(),
"data/".to_string(),
store,
polars::prelude::cloud::CloudOptions::default(),
&crate::OpenOptions::default(),
rt.handle(),
&crate::loading::measurements::OpenReport {
progress: Arc::new(crate::formats::schema_union::FooterProgress::default()),
meter: meter.clone(),
remembered: Some(cache.clone()),
writes: Default::default(),
},
);
meter
};
let first = open(store.clone());
assert_eq!(
first.footers().and_then(|c| c.files),
Some(2),
"the first open reads both footers"
);
let second = open(store.clone());
assert_eq!(
second.listing().and_then(|c| c.files),
Some(2),
"the second lists the prefix, which is how it knows nothing has changed"
);
assert_eq!(
second.footers(),
None,
"and reads no footers at all, which is what remembering them is for"
);
rt.block_on(async {
store
.put(&OsPath::from("data/b.parquet"), PutPayload::from(body(9)))
.await
.unwrap();
});
let third = open(store);
assert_eq!(
third.footers().and_then(|c| c.files),
Some(2),
"a dataset that has changed is read again rather than remembered wrongly"
);
}
#[test]
fn a_glob_reaches_the_route_that_lists_and_reads_it() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let body = || {
let mut frame = df!("n" => &[1i64]).unwrap();
let mut out = Vec::new();
ParquetWriter::new(&mut out).finish(&mut frame).unwrap();
out
};
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
rt.block_on(async {
for key in [
"data/year=2024/a.parquet",
"data/year=2025/b.parquet",
"data/other/c.parquet",
] {
store
.put(&OsPath::from(key), PutPayload::from(body()))
.await
.unwrap();
}
});
let meter = Arc::new(crate::loading::measurements::Meter::default());
let _ = crate::App::schema_state_from_cloud_hive_with(
"memory://data/year=*/*.parquet".to_string(),
"data/year=*/*.parquet".to_string(),
store,
polars::prelude::cloud::CloudOptions::default(),
&crate::OpenOptions::default(),
rt.handle(),
&crate::loading::measurements::OpenReport {
progress: Arc::new(crate::formats::schema_union::FooterProgress::default()),
meter: meter.clone(),
remembered: None,
writes: Default::default(),
},
);
assert_eq!(
meter.listing().and_then(|c| c.files),
Some(2),
"the listing found the two files the glob names — not the sibling directory \
it does not. A starred key handed to the listing finds none of them"
);
assert_eq!(
meter.footers().and_then(|c| c.files),
Some(2),
"and their footers were read, which is what a glob used to get none of"
);
let full = "s3://bucket/data/year=*/*.parquet";
let root = url_of_key(full, prefix_of_glob("data/year=*/*.parquet")).unwrap();
assert_eq!(root, "s3://bucket/data");
let file_url = url_of_key(full, "data/year=2024/a.parquet").unwrap();
assert!(
file_url.starts_with(&root),
"{file_url} has to sit under {root}, or every note measured from the root \
is silently about no files at all"
);
assert!(!file_url.starts_with(full), "which the URL as typed is not");
}
#[test]
fn a_glob_opens_the_files_it_names_and_no_others() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let body = || {
let mut frame = df!("n" => &[1i64]).unwrap();
let mut out = Vec::new();
ParquetWriter::new(&mut out).finish(&mut frame).unwrap();
out
};
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
rt.block_on(async {
for key in [
"data/year=2024/a.parquet",
"data/year=2025/b.parquet",
"data/other/c.parquet",
"elsewhere/d.parquet",
] {
store
.put(&OsPath::from(key), PutPayload::from(body()))
.await
.unwrap();
}
});
assert_eq!(prefix_of_glob("data/year=*/*.parquet"), "data");
assert_eq!(prefix_of_glob("data/*.parquet"), "data");
assert_eq!(prefix_of_glob("*.parquet"), "");
assert_eq!(prefix_of_glob("data/plain.parquet"), "data/plain.parquet");
let matcher = globset::GlobBuilder::new("data/year=*/*.parquet")
.literal_separator(true)
.build()
.unwrap()
.compile_matcher();
let (files, _skipped) = rt
.block_on(async {
list_dataset_files(
&store,
prefix_of_glob("data/year=*/*.parquet"),
Some(&matcher),
)
.await
})
.expect("the prefix lists");
let keys: Vec<&str> = files.iter().map(|f| f.key.as_str()).collect();
assert_eq!(
keys,
["data/year=2024/a.parquet", "data/year=2025/b.parquet"],
"the two the glob names — not the sibling directory it does not, and not the \
one outside the prefix altogether"
);
assert!(
!matcher.is_match("data/year=2024/deeper/a.parquet"),
"a single star does not cross a slash"
);
}
#[test]
fn a_directory_is_the_same_table_from_a_disk_or_a_bucket() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let body = || {
let mut frame = df!("n" => &[1i64]).unwrap();
let mut out = Vec::new();
ParquetWriter::new(&mut out).finish(&mut frame).unwrap();
out
};
let layout = [
("date=1/data.parquet", true),
("date=1/_2024.parquet", false),
("occurrence.parquet/part-00001", true),
("date=1/notes.csv", false),
];
let dir = tempfile::tempdir().unwrap();
for (rel, _) in layout {
let path = dir.path().join(rel);
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
std::fs::write(&path, body()).unwrap();
}
let (local, _skipped) = crate::formats::dataset_files::LocalFiles::new(dir.path()).walk(None);
let mut from_disk: Vec<String> = local
.iter()
.map(|p| {
p.strip_prefix(dir.path())
.unwrap()
.to_string_lossy()
.replace('\\', "/")
})
.collect();
from_disk.sort();
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
rt.block_on(async {
for (rel, _) in layout {
store
.put(
&OsPath::from(format!("data/{rel}")),
PutPayload::from(body()),
)
.await
.unwrap();
}
});
let (cloud, _skipped) = rt
.block_on(async { list_dataset_files(&store, "data/", None).await })
.expect("the prefix lists");
let mut from_bucket: Vec<String> = cloud
.iter()
.map(|f| f.key.trim_start_matches("data/").to_string())
.collect();
from_bucket.sort();
let mut wanted: Vec<String> = layout
.iter()
.filter(|(_, keep)| *keep)
.map(|(rel, _)| rel.to_string())
.collect();
wanted.sort();
assert_eq!(from_disk, wanted, "the disk reads the table");
assert_eq!(from_bucket, wanted, "and the bucket reads the same one");
}
#[test]
fn partitions_are_the_same_from_a_disk_or_a_bucket() {
use object_store::PutPayload;
use polars::prelude::{DataType, ParquetWriter, df};
let body = || {
let mut frame = df!("n" => &[1i64]).unwrap();
let mut out = Vec::new();
ParquetWriter::new(&mut out).finish(&mut frame).unwrap();
out
};
let layout = [
"k=1/a.parquet",
"k=2x/b.parquet",
"k=3/m=2024-01-01/c.parquet",
];
let partitions = |state: &crate::table::DataTableState| {
let columns = state.partition_columns().unwrap_or_default().to_vec();
columns
.iter()
.map(|c| (c.clone(), state.schema().get(c).cloned()))
.collect::<Vec<_>>()
};
let dir = tempfile::tempdir().unwrap();
for rel in layout {
let path = dir.path().join(rel);
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
std::fs::write(&path, body()).unwrap();
}
let report = || crate::loading::measurements::OpenReport {
progress: Arc::new(crate::formats::schema_union::FooterProgress::default()),
meter: Arc::new(crate::loading::measurements::Meter::default()),
remembered: None,
writes: Default::default(),
};
let (local, _) = crate::App::schema_state_from_local_hive(
Some(dir.path()),
&crate::OpenOptions {
hive: true,
..crate::OpenOptions::default()
},
&report(),
)
.expect("the directory opens");
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
rt.block_on(async {
for rel in layout {
store
.put(
&OsPath::from(format!("data/{rel}")),
PutPayload::from(body()),
)
.await
.unwrap();
}
});
let (cloud, _) = crate::App::schema_state_from_cloud_hive_with(
"memory://data/".to_string(),
"data/".to_string(),
store,
polars::prelude::cloud::CloudOptions::default(),
&crate::OpenOptions::default(),
rt.handle(),
&report(),
)
.expect("the prefix opens");
let from_disk = partitions(&local);
assert_eq!(
from_disk,
[
("k".to_string(), Some(DataType::Int64)),
("m".to_string(), Some(DataType::Date)),
],
"the newest file names the columns and the two ends type them"
);
assert_eq!(partitions(&cloud), from_disk, "and the bucket agrees");
}
#[test]
fn a_cloud_open_measures_what_its_listing_and_its_footers_cost() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let write = |rows: i64| -> Vec<u8> {
let mut frame = df!("n" => (0..rows).collect::<Vec<i64>>()).unwrap();
let mut bytes = Vec::new();
ParquetWriter::new(&mut bytes).finish(&mut frame).unwrap();
bytes
};
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
let mut written = 0u64;
rt.block_on(async {
for (key, rows) in [
("data/date=2024-01-01/a.parquet", 1),
("data/date=2024-01-02/b.parquet", 2),
("data/date=2024-01-03/c.parquet", 3),
] {
let body = write(rows);
written += body.len() as u64;
store
.put(&OsPath::from(key), PutPayload::from(body))
.await
.unwrap();
}
});
let meter = Arc::new(crate::loading::measurements::Meter::default());
let _ = crate::App::schema_state_from_cloud_hive_with(
"memory://data/".to_string(),
"data/".to_string(),
store,
polars::prelude::cloud::CloudOptions::default(),
&crate::OpenOptions::default(),
rt.handle(),
&crate::loading::measurements::OpenReport {
progress: Arc::new(crate::formats::schema_union::FooterProgress::default()),
meter: meter.clone(),
remembered: None,
writes: Default::default(),
},
);
let listing = meter.listing().expect("the open measured its listing");
assert_eq!(listing.files, Some(3), "the listing returned three objects");
assert!(
listing.over_the_wire.is_none(),
"object_store turns the listing's pages over itself, so the requests are \
not datui's to count and it claims none"
);
let footers = meter.footers().expect("the open measured its footer pass");
assert_eq!(
footers.files,
Some(3),
"a footer was read from each of the three"
);
let wire = footers
.over_the_wire
.expect("datui issued the footer reads itself, so it counts them");
assert_eq!(
wire.requests, 3,
"one ranged read each: these footers fit inside the tail datui asks for"
);
assert!(
wire.bytes.is_some_and(|b| b > 0 && b <= written),
"the bytes are what those reads returned — some, and no more than the three \
objects hold ({} of {written})",
wire.bytes.unwrap_or(0)
);
let total = meter.total().expect("and a total over both");
assert_eq!(
total.over_the_wire.map(|w| w.requests),
Some(3),
"the total carries the requests of the stretch that made any"
);
}
#[test]
fn a_column_only_a_middle_file_has_joins_after_the_open() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let plain = |i: i64| -> Vec<u8> {
let mut frame = df!("id" => &[i], "v" => &[i * 2]).unwrap();
let mut bytes = Vec::new();
ParquetWriter::new(&mut bytes).finish(&mut frame).unwrap();
bytes
};
let with_oops = |i: i64| -> Vec<u8> {
let mut frame = df!("id" => &[i], "v" => &[i * 2], "oops" => &["vendor"]).unwrap();
let mut bytes = Vec::new();
ParquetWriter::new(&mut bytes).finish(&mut frame).unwrap();
bytes
};
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
let files = FOOTERS_AT_ONCE + 1;
let odd_one_out = files / 2;
rt.block_on(async {
for i in 0..files {
let key = format!("data/date=2024-01-{:03}/part.parquet", i + 1);
let body = if i == odd_one_out {
with_oops(i as i64)
} else {
plain(i as i64)
};
store
.put(&OsPath::from(key.as_str()), PutPayload::from(body))
.await
.unwrap();
}
});
let progress = Arc::new(crate::formats::schema_union::FooterProgress::default());
let mut state = crate::App::schema_state_from_cloud_hive_with(
"memory://data/".to_string(),
"data/".to_string(),
store,
polars::prelude::cloud::CloudOptions::default(),
&crate::OpenOptions::default(),
rt.handle(),
&crate::loading::measurements::OpenReport {
progress: progress.clone(),
meter: Arc::new(crate::loading::measurements::Meter::default()),
remembered: None,
writes: Default::default(),
},
)
.map(|(state, facts)| {
state.with_open(crate::table::OpenFacts {
remote_source: true,
..facts
})
})
.expect("the prefix opens");
assert_eq!(
progress.last_pass().read,
2,
"the open waited for two footers, not {files}"
);
assert!(
!state.get_column_order().iter().any(|c| c == "oops"),
"the two ends cannot know about a column only the middle has: {:?}",
state.get_column_order()
);
assert!(
state.remote_files_counter().is_none(),
"the pass already reading every footer is where the count comes from"
);
state.scroll_right();
let scrolled_to = state.termcol_index;
assert!(scrolled_to > 0, "the fixture can be scrolled");
let join = state
.footers_pending()
.expect("the rest are still to be read");
let found = join(&progress).expect("the pass reads them");
state.visible_rows = 10;
assert!(
state.join_dataset_schema(found).is_ok(),
"nothing is built on top of the scan here, so they go straight in"
);
assert_eq!(
progress.last_pass().read,
files - 2,
"the pass behind the open read every footer but the two the open did"
);
assert_eq!(
state.get_column_order().last().map(String::as_str),
Some("oops"),
"the column joins, at the end, where nothing already shown has to move: \
{:?}",
state.get_column_order()
);
assert_eq!(
state.termcol_index, scrolled_to,
"and the view does not move under the user to make room"
);
assert!(
state.footers_pending().is_none(),
"with nothing left to wait for"
);
let mut request = state
.prepare_async_collect(None)
.expect("a page is planned");
let planned = request.lf.collect_schema();
assert!(
planned.is_ok(),
"the first page after the join could not even be planned: {:?}",
planned.err()
);
let planned = planned.unwrap();
assert!(
planned.iter_names().any(|name| name == "oops"),
"and it reads the column that just joined: {:?}",
planned.iter_names().collect::<Vec<_>>()
);
assert_eq!(
state.num_rows_if_valid(),
Some(files),
"and the count the pass brought back with it, one row a file — without a \
second pass over the same footers to learn it"
);
assert!(
state.display_df().is_none(),
"the frame is rebuilt but not read; the caller reads it back off the loop"
);
}
#[test]
fn an_object_only_the_pass_finds_corrupt_is_left_out_by_the_pass() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let body = |i: i64| -> Vec<u8> {
let mut frame = df!("id" => &[i]).unwrap();
let mut bytes = Vec::new();
ParquetWriter::new(&mut bytes).finish(&mut frame).unwrap();
bytes
};
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
let files = FOOTERS_AT_ONCE + 1;
let unreadable = files / 2;
rt.block_on(async {
for i in 0..files {
let key = format!("data/date=2024-01-{:03}/part.parquet", i + 1);
let bytes = if i == unreadable {
b"not a parquet file".to_vec()
} else {
body(i as i64)
};
store
.put(&OsPath::from(key.as_str()), PutPayload::from(bytes))
.await
.unwrap();
}
});
let progress = Arc::new(crate::formats::schema_union::FooterProgress::default());
let mut state = crate::App::schema_state_from_cloud_hive_with(
"memory://data/".to_string(),
"data/".to_string(),
store,
polars::prelude::cloud::CloudOptions::default(),
&crate::OpenOptions::default(),
rt.handle(),
&crate::loading::measurements::OpenReport {
progress: progress.clone(),
meter: Arc::new(crate::loading::measurements::Meter::default()),
remembered: None,
writes: Default::default(),
},
)
.map(|(state, facts)| state.with_open(facts))
.expect("the prefix opens");
let join = state
.footers_pending()
.expect("the rest are still to be read");
let found = join(&progress).expect("the pass reads them");
assert!(
state.join_dataset_schema(found).is_ok(),
"nothing is built on top of the scan here"
);
let plan = state
.visible_lf()
.explain(false)
.expect("the scan can be planned");
assert!(
!plan.contains(&format!("date=2024-01-{:03}", unreadable + 1)),
"the object that will not parse is not one of the sources: {plan}"
);
let counter = state
.remote_files_counter()
.expect("the dataset has not counted itself yet");
let groups = counter(&Default::default()).expect("the readable objects are counted");
let total = groups.iter().flatten().sum();
assert!(state.count_landed(state.len_generation(), total, Some(&groups)));
assert_eq!(
state.num_rows_if_valid(),
Some(files - 1),
"and the count lands — one row from every object that would open, rather \
than an answer for a list the dataset no longer holds, dropped in silence"
);
}
#[test]
fn a_dataset_of_one_wave_of_footers_opens_whole() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let body = |i: i64| -> Vec<u8> {
let mut frame = df!("id" => &[i]).unwrap();
let mut bytes = Vec::new();
ParquetWriter::new(&mut bytes).finish(&mut frame).unwrap();
bytes
};
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
rt.block_on(async {
for i in 0..FOOTERS_AT_ONCE {
let key = format!("data/date=2024-01-{:03}/part.parquet", i + 1);
store
.put(
&OsPath::from(key.as_str()),
PutPayload::from(body(i as i64)),
)
.await
.unwrap();
}
});
let progress = Arc::new(crate::formats::schema_union::FooterProgress::default());
let state = crate::App::schema_state_from_cloud_hive_with(
"memory://data/".to_string(),
"data/".to_string(),
store,
polars::prelude::cloud::CloudOptions::default(),
&crate::OpenOptions::default(),
rt.handle(),
&crate::loading::measurements::OpenReport {
progress: progress.clone(),
meter: Arc::new(crate::loading::measurements::Meter::default()),
remembered: None,
writes: Default::default(),
},
)
.map(|(state, facts)| state.with_open(facts))
.expect("the prefix opens");
assert_eq!(
progress.last_pass().read,
FOOTERS_AT_ONCE,
"a wave's worth is read at the open, not two of them"
);
assert!(
state.footers_pending().is_none(),
"with nothing left to read behind it"
);
assert_eq!(
state.num_rows_if_valid(),
Some(FOOTERS_AT_ONCE),
"counted from those footers as it opens, rather than left to a later pass"
);
}
#[test]
fn a_dataset_is_listed_once_and_counted_from_its_footers() {
use object_store::PutPayload;
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
rt.block_on(async {
for (key, bytes) in evolving_dataset() {
store
.put(&OsPath::from(key.as_str()), PutPayload::from(bytes))
.await
.unwrap();
}
store
.put(
&OsPath::from("elsewhere/x.parquet"),
PutPayload::from(vec![1u8]),
)
.await
.unwrap();
let (files, _skipped) = list_dataset_files(&store, "data/", None).await.unwrap();
let keys: Vec<&str> = files.iter().map(|f| f.key.as_str()).collect();
assert_eq!(
keys,
[
"data/date=2009-01-03/part-0.parquet",
"data/date=2026-09-17/part-0.parquet"
]
);
let footers = footers_of_files(
&store,
&files,
&[0, 1],
&Arc::new(crate::loading::measurements::Meter::default()),
)
.await;
let groups: Vec<Vec<usize>> = footers
.into_iter()
.map(|f| f.map(|f| f.row_group_rows).unwrap_or_default())
.collect();
assert_eq!(groups, [vec![2], vec![5]]);
let (dataset, partitions) = schema_of(&store, &files).await;
let schema = dataset.schema;
assert_eq!(partitions, ["date"]);
let names: Vec<&str> = schema.iter_names().map(|n| n.as_str()).collect();
assert_eq!(
names,
["date", "id", "fee", "input"],
"the newest file's columns"
);
assert!(
format!("{:?}", schema.get("input").unwrap()).contains("address"),
"and its struct fields"
);
});
}
#[test]
fn a_lenient_scan_reads_files_written_years_apart() {
let dir = tempfile::TempDir::new().unwrap();
let mut urls = Vec::new();
for (key, bytes) in evolving_dataset().into_iter().take(2) {
let path = dir.path().join(&key);
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
std::fs::write(&path, bytes).unwrap();
urls.push(path.to_string_lossy().into_owned());
}
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> =
Arc::new(object_store::local::LocalFileSystem::new_with_prefix(dir.path()).unwrap());
let schema = rt.block_on(async {
let (files, _skipped) = list_dataset_files(&store, "data", None).await.unwrap();
schema_of(&store, &files).await.0.schema
});
let df = lenient_scan(&urls, schema, None, None, &[])
.unwrap()
.collect()
.unwrap();
assert_eq!(df.height(), 7);
let fees = df.column("fee").unwrap();
assert_eq!(fees.null_count(), 2, "the old file has no fee");
let second_file = lenient_scan(&urls[1..], df.schema().clone(), None, None, &[])
.unwrap()
.slice(3, 2)
.collect()
.unwrap();
assert_eq!(
second_file
.column("id")
.unwrap()
.i64()
.unwrap()
.into_no_null_iter()
.collect::<Vec<_>>(),
[6, 7]
);
}
fn open_dataset(
files: Vec<(String, Vec<u8>)>,
) -> (
crate::formats::schema_union::DatasetSchema,
polars::prelude::DataFrame,
tempfile::TempDir,
) {
let dir = tempfile::TempDir::new().unwrap();
let mut urls = Vec::new();
for (key, bytes) in &files {
let path = dir.path().join(key);
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
std::fs::write(&path, bytes).unwrap();
urls.push(path.to_string_lossy().into_owned());
}
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> =
Arc::new(object_store::local::LocalFileSystem::new_with_prefix(dir.path()).unwrap());
let (dataset, listed, file_rows) = rt.block_on(async {
let (listed, _skipped) = list_dataset_files(&store, "data", None).await.unwrap();
let read: Vec<usize> = (0..listed.len()).collect();
let footers = footers_of_files(
&store,
&listed,
&read,
&Arc::new(crate::loading::measurements::Meter::default()),
)
.await;
let rows: Vec<usize> = footers
.iter()
.map(|f| {
f.as_ref()
.map(|f| f.row_group_rows.iter().sum())
.unwrap_or(0)
})
.collect();
(schema_of(&store, &listed).await.0, listed, rows)
});
let urls: Vec<String> = listed
.iter()
.map(|f| dir.path().join(&f.key).to_string_lossy().into_owned())
.collect();
let drift = crate::formats::schema_union::ScanDrift::new(&urls, &dataset, &file_rows);
let mut df = lenient_scan(&urls, dataset.schema.clone(), None, drift.as_ref(), &[])
.unwrap()
.collect()
.unwrap();
let _ = df.drop_in_place(crate::formats::schema_union::DRIFT_COLUMN);
(dataset, df, dir)
}
fn parquet(df: polars::prelude::DataFrame) -> Vec<u8> {
use polars::prelude::ParquetWriter;
let mut df = df;
let mut bytes = Vec::new();
ParquetWriter::new(&mut bytes).finish(&mut df).unwrap();
bytes
}
#[test]
fn a_column_only_a_middle_file_has_is_not_hidden() {
use polars::prelude::df;
let files = vec![
(
"data/a.parquet".to_string(),
parquet(df!("id" => &[1i64]).unwrap()),
),
(
"data/b.parquet".to_string(),
parquet(df!("id" => &[2i64], "oops" => &["x"]).unwrap()),
),
(
"data/c.parquet".to_string(),
parquet(df!("id" => &[3i64]).unwrap()),
),
];
let (dataset, df, _dir) = open_dataset(files);
let names: Vec<&str> = dataset.schema.iter_names().map(|n| n.as_str()).collect();
assert_eq!(names, ["id", "oops"]);
assert_eq!(df.height(), 3);
assert_eq!(df.column("oops").unwrap().null_count(), 2);
}
#[test]
fn files_of_different_integer_widths_open_as_the_wider_one() {
use polars::prelude::{DataType, df};
let files = vec![
(
"data/a.parquet".to_string(),
parquet(df!("n" => &[1i32, 2]).unwrap()),
),
(
"data/b.parquet".to_string(),
parquet(df!("n" => &[3i64]).unwrap()),
),
];
let (dataset, df, _dir) = open_dataset(files);
assert_eq!(dataset.schema.get("n"), Some(&DataType::Int64));
assert_eq!(df.height(), 3);
assert_eq!(df.column("n").unwrap().null_count(), 0);
}
#[test]
fn a_number_and_text_column_keeps_the_rows_of_both() {
use polars::prelude::{DataType, df};
let files = vec![
(
"data/a.parquet".to_string(),
parquet(df!("price" => &["1", "2"]).unwrap()),
),
(
"data/b.parquet".to_string(),
parquet(df!("price" => &[3i64, 4, 5]).unwrap()),
),
];
let (dataset, df, _dir) = open_dataset(files);
assert_eq!(
dataset.schema.get("price"),
Some(&DataType::Int64),
"the type most rows have"
);
assert_eq!(df.height(), 5, "every row is still there");
assert_eq!(
df.column("price").unwrap().null_count(),
2,
"the text file is not read for the column"
);
let drifting: Vec<_> = dataset.drifting().map(|c| c.name.to_string()).collect();
assert_eq!(drifting, ["price"]);
assert_eq!(dataset.columns[0].conflicting_types, [DataType::String]);
}
#[test]
fn one_corrupt_file_does_not_stop_the_dataset_opening() {
use polars::prelude::df;
let files = vec![
(
"data/a.parquet".to_string(),
parquet(df!("id" => &[1i64]).unwrap()),
),
("data/b.parquet".to_string(), b"not a parquet file".to_vec()),
(
"data/c.parquet".to_string(),
parquet(df!("id" => &[3i64]).unwrap()),
),
];
let dir = tempfile::TempDir::new().unwrap();
for (key, bytes) in &files {
let path = dir.path().join(key);
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
std::fs::write(&path, bytes).unwrap();
}
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> =
Arc::new(object_store::local::LocalFileSystem::new_with_prefix(dir.path()).unwrap());
let (dataset, listed) = rt.block_on(async {
let (listed, _skipped) = list_dataset_files(&store, "data", None).await.unwrap();
(schema_of(&store, &listed).await.0, listed)
});
assert_eq!(dataset.unreadable, [1], "named, and left out of the scan");
let names: Vec<&str> = dataset.schema.iter_names().map(|n| n.as_str()).collect();
assert_eq!(names, ["id"]);
assert_eq!(listed.len(), 3);
let paths: Vec<String> = files
.iter()
.map(|(key, _)| dir.path().join(key).to_string_lossy().into_owned())
.collect();
let readable = crate::formats::schema_union::readable_paths(&paths, &dataset.unreadable);
assert_eq!(
readable.len(),
2,
"the one that will not parse is not scanned"
);
let rows = crate::formats::schema_union::lenient_scan(
&readable,
dataset.schema.clone(),
None,
None,
&[],
)
.and_then(|lf| lf.collect());
assert_eq!(
rows.map(|df| df.height()).ok(),
Some(2),
"the two readable files' rows"
);
let all =
crate::formats::schema_union::lenient_scan(&paths, dataset.schema.clone(), None, None, &[])
.and_then(|lf| lf.collect());
assert!(
all.is_err(),
"left in, it takes the readable files down with it"
);
}
#[test]
fn the_count_reads_only_what_the_open_did_not_and_remembers_the_dataset() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
let files = rt.block_on(async {
for day in 1..=5i64 {
let mut frame = df!("id" => (0..day).collect::<Vec<i64>>()).unwrap();
let mut bytes = Vec::new();
ParquetWriter::new(&mut bytes).finish(&mut frame).unwrap();
store
.put(
&OsPath::from(format!("data/date=2024-01-0{day}/a.parquet")),
PutPayload::from(bytes),
)
.await
.unwrap();
}
list_dataset_files(&store, "data/", None).await.unwrap().0
});
let files = Arc::new(files);
let full = "memory://data/";
let fingerprint = crate::cache::DatasetShape::fingerprint_of(
files
.iter()
.map(|f| (f.key.as_str(), f.size, f.stamp, f.etag.as_deref())),
);
let dir = tempfile::tempdir().unwrap();
let cache = crate::cache::CacheManager::with_dir(dir.path().to_path_buf());
let meter = Arc::new(crate::loading::measurements::Meter::default());
meter.listed(std::time::Duration::from_millis(1), Some(5), false);
let sampled = rt.block_on(footers_of_files(
&store,
&files,
&[0, 4],
&Arc::new(crate::loading::measurements::Meter::default()),
));
let mut planted = sampled[0].clone().unwrap();
planted.row_group_rows = vec![999];
let source = Arc::new(crate::formats::dataset_files::StoreFiles::new(
full,
"data/".to_string(),
None,
store.clone(),
polars::prelude::cloud::CloudOptions::default(),
rt.handle(),
));
let count = crate::formats::dataset_files::counter_for(
source,
files.to_vec(),
[(0, Some(planted)), (4, sampled[1].clone())],
Some(fingerprint.clone()),
meter.clone(),
Some(cache.clone()),
);
let groups = count(&Default::default()).unwrap();
assert_eq!(groups, [vec![999], vec![2], vec![3], vec![4], vec![5]]);
assert_eq!(
meter.footers().and_then(|c| c.files),
Some(3),
"the count read the three footers the open had not"
);
assert_eq!(
meter
.footers()
.and_then(|c| c.over_the_wire)
.map(|w| w.requests),
Some(3),
"one request each, and none for the two it was given"
);
let shape = cache
.dataset_shape(full, &fingerprint)
.expect("every footer is in, so the dataset is remembered");
assert_eq!(shape.files.len(), 5);
let progress = Arc::new(crate::formats::schema_union::FooterProgress::default());
let reopen = Arc::new(crate::loading::measurements::Meter::default());
let (state, facts) = crate::App::schema_state_from_cloud_hive_with(
full.to_string(),
"data/".to_string(),
store,
polars::prelude::cloud::CloudOptions::default(),
&crate::OpenOptions::default(),
rt.handle(),
&crate::loading::measurements::OpenReport {
progress: progress.clone(),
meter: reopen.clone(),
remembered: Some(cache),
writes: Default::default(),
},
)
.expect("the reopen opens");
assert_eq!(progress.last_pass().begun, 0, "no footer pass");
assert!(reopen.footers().is_none(), "and no footer read");
assert!(facts.footers_pending.is_none(), "nor one behind the open");
let state = state.with_open(facts);
assert_eq!(state.num_rows_if_valid(), Some(999 + 2 + 3 + 4 + 5));
}
#[test]
fn a_dataset_with_a_corrupt_object_still_counts_the_rest() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let body = |i: i64| -> Vec<u8> {
let mut frame = df!("id" => &[i]).unwrap();
let mut bytes = Vec::new();
ParquetWriter::new(&mut bytes).finish(&mut frame).unwrap();
bytes
};
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
rt.block_on(async {
for (key, bytes) in [
("data/date=2024-01-01/a.parquet", body(1)),
(
"data/date=2024-01-02/b.parquet",
b"not a parquet file".to_vec(),
),
("data/date=2024-01-03/c.parquet", body(3)),
] {
store
.put(&OsPath::from(key), PutPayload::from(bytes))
.await
.unwrap();
}
});
let progress = Arc::new(crate::formats::schema_union::FooterProgress::default());
let mut state = crate::App::schema_state_from_cloud_hive_with(
"memory://data/".to_string(),
"data/".to_string(),
store,
polars::prelude::cloud::CloudOptions::default(),
&crate::OpenOptions::default(),
rt.handle(),
&crate::loading::measurements::OpenReport {
progress: progress.clone(),
meter: Arc::new(crate::loading::measurements::Meter::default()),
remembered: None,
writes: Default::default(),
},
)
.map(|(state, facts)| state.with_open(facts))
.expect("the prefix opens despite the one that will not parse");
let plan = state
.visible_lf()
.explain(false)
.expect("the scan can be planned");
assert!(
!plan.contains("b.parquet"),
"the object that will not parse is not one of the sources: {plan}"
);
assert!(
plan.contains("a.parquet") && plan.contains("c.parquet"),
"and the two that will are: {plan}"
);
let counter = state
.remote_files_counter()
.expect("the dataset has not counted itself yet, so it offers to");
let groups = counter(&Default::default()).expect("the readable objects are counted");
let total = groups.iter().flatten().sum();
assert!(state.count_landed(state.len_generation(), total, Some(&groups)));
assert_eq!(
state.num_rows_if_valid(),
Some(2),
"one row from each object that would open, and the count lands rather than \
being dropped on a length nobody mentions"
);
}
#[test]
fn a_directory_with_no_data_in_it_is_nobodys_table() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let parquet = || -> Vec<u8> {
let mut frame = df!("id" => &[1i64]).unwrap();
let mut bytes = Vec::new();
ParquetWriter::new(&mut bytes).finish(&mut frame).unwrap();
bytes
};
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
rt.block_on(async {
for (key, body) in [
("t/data/date=1/part-0.parquet", parquet()),
("t/metadata/v1.metadata.json", b"{}".to_vec()),
("t/metadata/v2.metadata.json", b"{}".to_vec()),
("t/metadata/snap-123.avro", b"x".to_vec()),
("t/metadata/version-hint.text", b"2".to_vec()),
("t/data/date=1", Vec::new()),
("t/data/date=2", Vec::new()),
("t/data/date=1/extra.csv", b"id\n1\n".to_vec()),
("t/data/date=1/part-1.parquet", Vec::new()),
("t/data/date=1/README", Vec::new()),
("t/data/y=2024/m=03/part-0.parquet", Vec::new()),
("t/data/date=4/part-0.csv", b"id\n4\n".to_vec()),
] {
store
.put(&OsPath::from(key), PutPayload::from(body))
.await
.unwrap();
}
});
let (files, skipped) = rt
.block_on(list_dataset_files(&store, "t", None))
.expect("the prefix lists");
assert_eq!(
files.iter().map(|f| f.key.as_str()).collect::<Vec<_>>(),
["t/data/date=1/part-0.parquet"]
);
assert_eq!(
skipped.not_parquet, 2,
"the csv beside the data and the day that landed as one: both are files \
somebody meant to be in the table"
);
assert_eq!(
skipped.empty, 2,
"and both whose names said Parquet over nothing at all — including the \
one alone in its partition, which is the case that matters most"
);
assert_eq!(
skipped.bookkeeping, 7,
"the four Iceberg files, the two folder markers, and a README with \
nothing in it: placeholders and plumbing"
);
}
#[test]
fn a_staged_open_does_not_lose_what_the_listing_passed_over() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let body = |i: i64| -> Vec<u8> {
let mut frame = df!("id" => &[i]).unwrap();
let mut bytes = Vec::new();
ParquetWriter::new(&mut bytes).finish(&mut frame).unwrap();
bytes
};
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
let files = FOOTERS_AT_ONCE + 1;
rt.block_on(async {
for i in 0..files {
let key = format!("data/date=2024-01-{:03}/part.parquet", i + 1);
store
.put(
&OsPath::from(key.as_str()),
PutPayload::from(body(i as i64)),
)
.await
.unwrap();
}
store
.put(
&OsPath::from("data/date=2024-01-001/extra.csv"),
PutPayload::from(b"id\n1\n".to_vec()),
)
.await
.unwrap();
});
let progress = Arc::new(crate::formats::schema_union::FooterProgress::default());
let mut state = crate::App::schema_state_from_cloud_hive_with(
"memory://data/".to_string(),
"data/".to_string(),
store,
polars::prelude::cloud::CloudOptions::default(),
&crate::OpenOptions::default(),
rt.handle(),
&crate::loading::measurements::OpenReport {
progress: progress.clone(),
meter: Arc::new(crate::loading::measurements::Meter::default()),
remembered: None,
writes: Default::default(),
},
)
.map(|(state, facts)| state.with_open(facts))
.expect("the prefix opens");
let said = |state: &crate::table::DataTableState| {
state
.notes()
.iter()
.any(|note| note.summary.contains("not Parquet"))
};
assert!(said(&state), "the open says so");
let join = state
.footers_pending()
.expect("the rest are still to be read");
let found = join(&progress).expect("the pass reads them");
assert!(state.join_dataset_schema(found).is_ok());
assert!(
said(&state),
"and it still does once the columns have joined: {:#?}",
state.notes()
);
}
#[test]
fn the_listing_counts_what_it_passes_over() {
use object_store::PutPayload;
use polars::prelude::{ParquetWriter, df};
let parquet = || -> Vec<u8> {
let mut frame = df!("id" => &[1i64]).unwrap();
let mut bytes = Vec::new();
ParquetWriter::new(&mut bytes).finish(&mut frame).unwrap();
bytes
};
let rt = tokio::runtime::Runtime::new().unwrap();
let store: Arc<dyn ObjectStore> = Arc::new(object_store::memory::InMemory::new());
rt.block_on(async {
for (key, body) in [
("data/date=1/part-0.parquet", parquet()),
("data/date=1/extra.csv", b"id\n1\n".to_vec()),
("data/date=1/notes.txt", b"read me".to_vec()),
("data/date=1/part-1.parquet", Vec::new()),
("data/_SUCCESS", Vec::new()),
("data/_delta_log/00000000000000000000.json", b"{}".to_vec()),
("data/_delta_log/.00000000000000000000.json.crc", Vec::new()),
] {
store
.put(&OsPath::from(key), PutPayload::from(body))
.await
.unwrap();
}
});
let (files, skipped) = rt
.block_on(list_dataset_files(&store, "data", None))
.expect("the prefix lists");
assert_eq!(
files.iter().map(|f| f.key.as_str()).collect::<Vec<_>>(),
["data/date=1/part-0.parquet"],
"one object is the table"
);
assert_eq!(
skipped,
crate::formats::schema_union::SkippedFiles {
not_parquet: 2,
empty: 1,
bookkeeping: 3,
},
"the csv and the txt are somebody's, the empty part is a write that \
stopped, and `_SUCCESS` and both log files are the writer's own"
);
}
#[test]
fn a_footer_from_a_tail_invalid_returns_err() {
let invalid = vec![0u8; 100];
let r = FileFooter::from_tail(&invalid, invalid.len(), true);
assert!(r.is_err());
}
#[test]
fn a_footer_from_a_tail_reads_schema_and_row_count() {
use polars::prelude::{ParquetWriter, df};
let mut df = df!("a" => &[1i32, 2, 3], "b" => &["x", "y", "z"]).unwrap();
let mut bytes = Vec::new();
ParquetWriter::new(&mut bytes).finish(&mut df).unwrap();
let footer = FileFooter::from_tail(&bytes, bytes.len(), true).unwrap();
assert_eq!(footer.row_group_rows, [3]);
assert_eq!(footer.schema.len(), 2);
}
#[test]
fn a_footer_from_a_tail_reads_row_groups_and_column_widths() {
use polars::prelude::{ParquetWriter, df};
let ids: Vec<i32> = (0..1000).collect();
let notes: Vec<String> = ids.iter().map(|i| format!("note-{i:04}")).collect();
let lists: Vec<Series> = ids
.iter()
.map(|i| Series::new("".into(), &[*i as f64; 10]))
.collect();
let mut df = df!("id" => ids, "note" => notes, "list" => lists).unwrap();
let mut bytes = Vec::new();
ParquetWriter::new(&mut bytes)
.with_row_group_size(Some(400))
.finish(&mut df)
.unwrap();
let footer = FileFooter::from_tail(&bytes, bytes.len(), true).unwrap();
assert!(
footer.row_group_rows.len() > 1,
"{:?}",
footer.row_group_rows
);
assert_eq!(footer.row_group_rows.iter().sum::<usize>(), 1000);
let widths = crate::formats::schema_union::column_bytes_per_row(&[Some(footer.clone())]);
let width = |column: &str| {
widths
.iter()
.find(|(n, _)| n == column)
.map(|(_, w)| *w)
.unwrap_or_else(|| panic!("no width for {column}"))
};
assert!(
(9..=20).contains(&width("note")),
"note width {}",
width("note")
);
assert!(
(80..=120).contains(&width("list")),
"list width {}",
width("list")
);
assert!(width("id") >= 4);
}