Read chunks of 4 GiB or more in every chunk index
ChunkInfo::chunk_size (and ChunkMapping::file_size) are u64: sizes past u32 were truncated for Single Chunk, Implicit, Fixed and Extensible Array indexes, and a v2 B-tree index refused them. A selection of a chunked dataset with a non-default fill value is read over a box of fill values instead of a full read, an unfiltered chunk of a file that is not in memory is read row by row, and an intermediate deflate stage no longer reserves the chunk's whole bound. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
This commit is contained in:
@@ -7,3 +7,6 @@ weights/
|
||||
.venv
|
||||
__pycache__/
|
||||
.pytest_cache/
|
||||
|
||||
# Scratch files the heavy tests generate (huge_chunks_interop)
|
||||
crates/*/tests/scratch/
|
||||
|
||||
@@ -899,7 +899,7 @@ mod tests {
|
||||
|
||||
fn make_chunk(offsets: Vec<u64>, address: u64, size: u32) -> ChunkInfo {
|
||||
ChunkInfo {
|
||||
chunk_size: size,
|
||||
chunk_size: u64::from(size),
|
||||
filter_mask: 0,
|
||||
offsets,
|
||||
address,
|
||||
|
||||
@@ -109,7 +109,7 @@ pub struct ChunkMapping {
|
||||
/// File byte address of the compressed chunk.
|
||||
pub file_offset: u64,
|
||||
/// Size of the compressed chunk in the file.
|
||||
pub file_size: u32,
|
||||
pub file_size: u64,
|
||||
/// Filter mask (0 = all filters applied).
|
||||
pub filter_mask: u32,
|
||||
/// Pre-computed row-copy operations for assembling this chunk into output.
|
||||
@@ -338,7 +338,7 @@ mod tests {
|
||||
|
||||
fn make_chunk(offsets: Vec<u64>, address: u64, size: u32) -> ChunkInfo {
|
||||
ChunkInfo {
|
||||
chunk_size: size,
|
||||
chunk_size: u64::from(size),
|
||||
filter_mask: 0,
|
||||
offsets,
|
||||
address,
|
||||
|
||||
@@ -314,7 +314,7 @@ pub(crate) fn chunk_req(
|
||||
chunk_bytes: usize,
|
||||
wanted: bool,
|
||||
) -> ExtentReq {
|
||||
let len = c.chunk_size as usize;
|
||||
let len = crate::addr::saturating_usize(c.chunk_size);
|
||||
ExtentReq {
|
||||
addr: c.address,
|
||||
len,
|
||||
@@ -521,7 +521,7 @@ pub fn decompress_all_chunks_with_stats_in<S: Storage + ?Sized>(
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct ChunkInfo {
|
||||
/// Size of chunk data in the file (after compression).
|
||||
pub chunk_size: u32,
|
||||
pub chunk_size: u64,
|
||||
/// Bitmask of filters that were NOT applied (0 = all applied).
|
||||
pub filter_mask: u32,
|
||||
/// N-dimensional offset of this chunk in dataset space.
|
||||
@@ -610,7 +610,7 @@ fn stored_element_size(dt: &Datatype, offset_size: u8) -> u64 {
|
||||
/// (`H5D__chunk_set_sizes`: "stored datatype size in chunk layout does not
|
||||
/// match datatype description"). Reading it anyway laid the chunks out with
|
||||
/// the wrong element size.
|
||||
pub(crate) fn check_chunk_element_size(
|
||||
pub fn check_chunk_element_size(
|
||||
layout: &DataLayout,
|
||||
datatype: &Datatype,
|
||||
offset_size: u8,
|
||||
@@ -1028,7 +1028,7 @@ fn parse_chunk_node<S: Storage + ?Sized>(
|
||||
if node_level == 0 {
|
||||
chunks.push(stored.len());
|
||||
stored.push(ChunkInfo {
|
||||
chunk_size,
|
||||
chunk_size: u64::from(chunk_size),
|
||||
filter_mask,
|
||||
offsets: keys[k..].to_vec(),
|
||||
address,
|
||||
@@ -1131,7 +1131,7 @@ pub fn generate_implicit_chunks_in_grid(
|
||||
}
|
||||
|
||||
chunks.push(ChunkInfo {
|
||||
chunk_size: chunk_byte_size as u32,
|
||||
chunk_size: chunk_byte_size,
|
||||
filter_mask: 0,
|
||||
offsets,
|
||||
address: base_address.saturating_add(grid_idx.saturating_mul(chunk_byte_size)),
|
||||
@@ -1190,9 +1190,15 @@ fn read_btree_v2_chunks<S: Storage + ?Sized>(
|
||||
}
|
||||
_ => return Err(bad("tree is not a chunk index")),
|
||||
};
|
||||
let unfiltered_bytes = checked_chunk_byte_len(chunk_dims, elem_size)?;
|
||||
let unfiltered_bytes =
|
||||
u32::try_from(unfiltered_bytes).map_err(|_| bad("chunk larger than 4 GiB"))?;
|
||||
// A u64 whatever the platform: an unfiltered chunk's size is only
|
||||
// recorded here; a chunk this platform cannot address fails when read.
|
||||
let unfiltered_bytes = u64::try_from(
|
||||
chunk_dims
|
||||
.iter()
|
||||
.try_fold(elem_size as u128, |acc, &c| acc.checked_mul(c as u128))
|
||||
.ok_or_else(|| bad("chunk size overflows"))?,
|
||||
)
|
||||
.map_err(|_| bad("chunk size overflows"))?;
|
||||
|
||||
let records = collect_btree_v2_records_in(file_data, &header, offset_size, length_size)?;
|
||||
let mut chunks = Vec::with_capacity(records.len());
|
||||
@@ -1213,10 +1219,7 @@ fn read_btree_v2_chunks<S: Storage + ?Sized>(
|
||||
pos += size_len;
|
||||
let mask = u32::from_le_bytes([data[pos], data[pos + 1], data[pos + 2], data[pos + 3]]);
|
||||
pos += 4;
|
||||
(
|
||||
u32::try_from(size).map_err(|_| bad("stored chunk larger than 4 GiB"))?,
|
||||
mask,
|
||||
)
|
||||
(size, mask)
|
||||
};
|
||||
let mut offsets = Vec::with_capacity(rank);
|
||||
for &dim in chunk_dims {
|
||||
@@ -1334,9 +1337,9 @@ pub fn list_chunks_in<S: Storage + ?Sized>(
|
||||
// Single chunk — one chunk covering the entire dataset
|
||||
let chunk_byte_size = checked_chunk_byte_len(&chunk_dims, elem_size)?;
|
||||
let (csize, fmask) = if let Some(fs) = single_filtered_size {
|
||||
(fs as u32, single_filter_mask.unwrap_or(0))
|
||||
(fs, single_filter_mask.unwrap_or(0))
|
||||
} else {
|
||||
(chunk_byte_size as u32, 0)
|
||||
(chunk_byte_size as u64, 0)
|
||||
};
|
||||
vec![ChunkInfo {
|
||||
chunk_size: csize,
|
||||
@@ -2124,7 +2127,7 @@ pub fn read_chunked_data_indexed_in<S: Storage + ?Sized>(
|
||||
.iter()
|
||||
.zip(&hits)
|
||||
.map(|(m, hit)| {
|
||||
let len = m.file_size as usize;
|
||||
let len = crate::addr::saturating_usize(m.file_size);
|
||||
ExtentReq {
|
||||
addr: m.file_offset,
|
||||
len,
|
||||
@@ -2757,7 +2760,7 @@ mod tests {
|
||||
}
|
||||
|
||||
chunk_infos.push(ChunkInfo {
|
||||
chunk_size: chunk_bytes as u32,
|
||||
chunk_size: chunk_bytes as u64,
|
||||
filter_mask: 0,
|
||||
offsets: vec![start as u64, 0],
|
||||
address: data_offset as u64,
|
||||
@@ -2869,7 +2872,7 @@ mod tests {
|
||||
.collect();
|
||||
let stored = crate::filters::compress_chunk(&chunk, &pipeline, 4).unwrap();
|
||||
chunks.push(ChunkInfo {
|
||||
chunk_size: stored.len() as u32,
|
||||
chunk_size: stored.len() as u64,
|
||||
filter_mask: 0,
|
||||
offsets: vec![r0 as u64, c0 as u64, 0],
|
||||
address: file.len() as u64,
|
||||
@@ -2943,7 +2946,7 @@ mod tests {
|
||||
let short = crate::filters::compress_chunk(&[1u8; 64], &pipeline, 4).unwrap();
|
||||
for bad in [5usize, 11, 40] {
|
||||
chunks[bad].address = file.len() as u64;
|
||||
chunks[bad].chunk_size = short.len() as u32;
|
||||
chunks[bad].chunk_size = short.len() as u64;
|
||||
file.extend_from_slice(&short);
|
||||
}
|
||||
for _ in 0..20 {
|
||||
@@ -3133,7 +3136,7 @@ mod tests {
|
||||
file_data[data_offset..data_offset + compressed.len()].copy_from_slice(&compressed);
|
||||
|
||||
chunk_infos.push(ChunkInfo {
|
||||
chunk_size: compressed.len() as u32,
|
||||
chunk_size: compressed.len() as u64,
|
||||
filter_mask: 0,
|
||||
offsets: vec![start as u64, 0],
|
||||
address: data_offset as u64,
|
||||
@@ -3215,7 +3218,7 @@ mod tests {
|
||||
file_data[data_offset..data_offset + chunk_size].copy_from_slice(&chunk_bytes);
|
||||
|
||||
chunk_infos.push(ChunkInfo {
|
||||
chunk_size: chunk_size as u32,
|
||||
chunk_size: chunk_size as u64,
|
||||
filter_mask: 0,
|
||||
offsets: vec![row_start as u64, col_start as u64, 0],
|
||||
address: data_offset as u64,
|
||||
@@ -3340,7 +3343,7 @@ mod tests {
|
||||
assert_eq!(c.address, 0x1000 + i as u64 * chunk_byte_size as u64);
|
||||
assert_eq!(c.offsets, vec![i as u64 * 20]);
|
||||
assert_eq!(c.filter_mask, 0);
|
||||
assert_eq!(c.chunk_size, chunk_byte_size as u32);
|
||||
assert_eq!(c.chunk_size, chunk_byte_size as u64);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -2125,7 +2125,7 @@ mod tests {
|
||||
for info in &infos {
|
||||
// Skipped chunks are stored at the chunk's size (shuffled).
|
||||
assert_eq!(
|
||||
info.chunk_size == (c * 8) as u32,
|
||||
info.chunk_size == (c * 8) as u64,
|
||||
info.filter_mask != 0,
|
||||
"{info:?}"
|
||||
);
|
||||
|
||||
@@ -224,7 +224,7 @@ fn read_element(
|
||||
};
|
||||
Ok((
|
||||
Some(ChunkInfo {
|
||||
chunk_size: chunk_byte_size as u32,
|
||||
chunk_size: chunk_byte_size,
|
||||
filter_mask: 0,
|
||||
offsets,
|
||||
address,
|
||||
@@ -259,7 +259,7 @@ fn read_element(
|
||||
};
|
||||
Ok((
|
||||
Some(ChunkInfo {
|
||||
chunk_size: chunk_size as u32,
|
||||
chunk_size,
|
||||
filter_mask,
|
||||
offsets,
|
||||
address,
|
||||
@@ -946,7 +946,7 @@ mod tests {
|
||||
assert_eq!(chunks.len(), 2);
|
||||
assert_eq!(chunks[0].address, base_addr);
|
||||
assert_eq!(chunks[0].offsets, vec![0]);
|
||||
assert_eq!(chunks[0].chunk_size, chunk_byte_size as u32);
|
||||
assert_eq!(chunks[0].chunk_size, chunk_byte_size as u64);
|
||||
assert_eq!(chunks[1].address, base_addr + chunk_byte_size);
|
||||
assert_eq!(chunks[1].offsets, vec![20]);
|
||||
}
|
||||
|
||||
@@ -340,10 +340,23 @@ pub fn decompress_chunk_exact_with<'s>(
|
||||
} else {
|
||||
MAX_DECOMPRESS_SIZE
|
||||
};
|
||||
let size_hint = if ctx.max_output != 0 {
|
||||
// The bound is the output's size when only size-preserving
|
||||
// filters (shuffle, Fletcher32) remain to be undone; before
|
||||
// any other filter (a second deflate, N-Bit, ...) it is only
|
||||
// a ceiling, and the output starts smaller and grows: the
|
||||
// inflater writes every byte it reserves, so reserving a
|
||||
// 4 GiB chunk's bound for a stage a few MiB long would hold
|
||||
// twice the chunk's memory.
|
||||
let exact = pipeline.filters[..i].iter().enumerate().all(|(j, f)| {
|
||||
filter_skipped(filter_mask, j)
|
||||
|| matches!(f.filter_id, FILTER_SHUFFLE | FILTER_FLETCHER32)
|
||||
});
|
||||
let size_hint = if ctx.max_output == 0 {
|
||||
input.len().saturating_mul(4).min(1 << 20)
|
||||
} else if exact {
|
||||
ctx.max_output
|
||||
} else {
|
||||
input.len().saturating_mul(4).min(1 << 20)
|
||||
input.len().saturating_mul(4).min(ctx.max_output)
|
||||
};
|
||||
let inflater = scratch
|
||||
.inflater
|
||||
|
||||
@@ -383,7 +383,7 @@ fn parse_fa_element(
|
||||
offset_size: u8,
|
||||
element_size: u8,
|
||||
chunk_byte_size: u64,
|
||||
) -> Result<Option<(u64, u32, u32)>, FormatError> {
|
||||
) -> Result<Option<(u64, u64, u32)>, FormatError> {
|
||||
let os = offset_size as usize;
|
||||
if client_id == 0 {
|
||||
// Non-filtered: element is just the chunk address.
|
||||
@@ -393,7 +393,7 @@ fn parse_fa_element(
|
||||
return Ok(None);
|
||||
}
|
||||
let address = read_offset(file_data, abs, offset_size)?;
|
||||
Ok(Some((address, chunk_byte_size as u32, 0)))
|
||||
Ok(Some((address, chunk_byte_size, 0)))
|
||||
} else {
|
||||
// Filtered: address(offset_size) + chunk_size(variable) + filter_mask(4)
|
||||
let es = element_size as usize;
|
||||
@@ -418,7 +418,7 @@ fn parse_fa_element(
|
||||
file_data[fm_off + 2],
|
||||
file_data[fm_off + 3],
|
||||
]);
|
||||
Ok(Some((address, chunk_size as u32, filter_mask)))
|
||||
Ok(Some((address, chunk_size, filter_mask)))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -688,7 +688,7 @@ mod tests {
|
||||
assert_eq!(c.address, base_addr + i as u64 * chunk_byte_size as u64);
|
||||
assert_eq!(c.offsets, vec![i as u64 * 20]);
|
||||
assert_eq!(c.filter_mask, 0);
|
||||
assert_eq!(c.chunk_size, chunk_byte_size as u32);
|
||||
assert_eq!(c.chunk_size, chunk_byte_size as u64);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -426,7 +426,7 @@ mod tests {
|
||||
for i in 0..8u64 {
|
||||
let len = if short && i == 5 { 16 } else { 32 };
|
||||
infos.push(ChunkInfo {
|
||||
chunk_size: len as u32,
|
||||
chunk_size: len as u64,
|
||||
filter_mask: 0,
|
||||
offsets: vec![i * 8],
|
||||
address: file.len() as u64,
|
||||
|
||||
@@ -243,6 +243,58 @@ fn copy_overlap(
|
||||
}
|
||||
}
|
||||
|
||||
/// Copy the part of unfiltered chunk `chunk` (`chunk_bytes` long, shape
|
||||
/// `chunk_shape`) that overlaps the box into `out`, reading only the runs of
|
||||
/// the overlap from the file. The whole chunk must still lie inside the
|
||||
/// file, as it must when it is fetched whole.
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
fn read_unfiltered_overlap<S: Storage + ?Sized>(
|
||||
file_data: &S,
|
||||
chunk: &crate::chunked_read::ChunkInfo,
|
||||
chunk_bytes: usize,
|
||||
chunk_shape: &[u64],
|
||||
elem_size: usize,
|
||||
out: &mut [u8],
|
||||
box_start: &[u64],
|
||||
box_extent: &[u64],
|
||||
) -> Result<(), FormatError> {
|
||||
let rank = chunk_shape.len();
|
||||
let origin = &chunk.offsets[..rank];
|
||||
let (lo, extent): (Vec<u64>, Vec<u64>) = (0..rank)
|
||||
.map(|d| {
|
||||
let lo = origin[d].max(box_start[d]);
|
||||
let hi = origin[d]
|
||||
.saturating_add(chunk_shape[d])
|
||||
.min(box_start[d] + box_extent[d]);
|
||||
(lo, hi.saturating_sub(lo))
|
||||
})
|
||||
.unzip();
|
||||
let file_len = crate::storage::len_usize(file_data);
|
||||
let base = crate::addr::to_usize(chunk.address)?;
|
||||
if base > file_len || chunk_bytes > file_len - base {
|
||||
return Err(FormatError::UnexpectedEof {
|
||||
expected: base.saturating_add(chunk_bytes),
|
||||
available: file_len,
|
||||
});
|
||||
}
|
||||
let overlap = Selection::Hyperslab {
|
||||
start: lo.iter().zip(origin).map(|(l, o)| l - o).collect(),
|
||||
stride: vec![1; rank],
|
||||
count: extent.clone(),
|
||||
block: vec![1; rank],
|
||||
};
|
||||
let rows = crate::gather::gather_storage(
|
||||
file_data,
|
||||
chunk.address,
|
||||
chunk_bytes,
|
||||
chunk_shape,
|
||||
elem_size,
|
||||
&overlap,
|
||||
)?;
|
||||
copy_overlap(&rows, &lo, &extent, out, box_start, box_extent, elem_size);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Read `selection` without materialising the whole dataset, when that is
|
||||
/// possible and worthwhile. `Ok(None)` means "use the full-read path": an
|
||||
/// `All`/`None`/invalid selection, a layout this doesn't handle (compact,
|
||||
@@ -281,6 +333,39 @@ pub fn read_selection_in<S: Storage + ?Sized>(
|
||||
offset_size: u8,
|
||||
length_size: u8,
|
||||
selection: &Selection,
|
||||
) -> Result<Option<Vec<u8>>, FormatError> {
|
||||
read_selection_filled_in(
|
||||
file_data,
|
||||
layout,
|
||||
dataspace,
|
||||
elem_size,
|
||||
pipeline,
|
||||
offset_size,
|
||||
length_size,
|
||||
selection,
|
||||
None,
|
||||
)
|
||||
}
|
||||
|
||||
/// [`read_selection_in`] for a dataset whose fill value is `fill` (one
|
||||
/// element's bytes; `None` or all zeros is the default fill): the elements
|
||||
/// of a chunked dataset's selection that lie in chunks never written read as
|
||||
/// `fill`, as they do in a full read. Only the chunks the selection's
|
||||
/// bounding box overlaps are read, so a selection of a few elements of a
|
||||
/// dataset whose chunks are 4 GiB or more costs one decoded chunk (a
|
||||
/// filtered chunk has to be decoded whole) or, unfiltered, only the bytes
|
||||
/// it selects.
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
pub fn read_selection_filled_in<S: Storage + ?Sized>(
|
||||
file_data: &S,
|
||||
layout: &DataLayout,
|
||||
dataspace: &Dataspace,
|
||||
elem_size: usize,
|
||||
pipeline: Option<&FilterPipeline>,
|
||||
offset_size: u8,
|
||||
length_size: u8,
|
||||
selection: &Selection,
|
||||
fill: Option<&[u8]>,
|
||||
) -> Result<Option<Vec<u8>>, FormatError> {
|
||||
let dims = &dataspace.dimensions;
|
||||
if dims.is_empty() || elem_size == 0 {
|
||||
@@ -341,8 +426,18 @@ pub fn read_selection_in<S: Storage + ?Sized>(
|
||||
return Ok(None);
|
||||
}
|
||||
let mut boxed = alloc_output(checked_byte_len(box_elements, elem_size)?)?;
|
||||
if let Some(fill) = fill.filter(|f| <[u8]>::len(f) == elem_size && f.iter().any(|&b| b != 0)) {
|
||||
for element in boxed.chunks_exact_mut(elem_size) {
|
||||
element.copy_from_slice(fill);
|
||||
}
|
||||
}
|
||||
|
||||
match layout {
|
||||
// No chunk was ever written: every element is the fill value.
|
||||
DataLayout::Chunked {
|
||||
btree_address: None,
|
||||
..
|
||||
} if fill.is_some() => {}
|
||||
DataLayout::Chunked {
|
||||
btree_address: Some(_),
|
||||
..
|
||||
@@ -373,6 +468,25 @@ pub fn read_selection_in<S: Storage + ?Sized>(
|
||||
})
|
||||
})
|
||||
.collect();
|
||||
// An unfiltered chunk of a file that is not in memory: fetch
|
||||
// only the rows the box needs, not the whole chunk (which may be
|
||||
// 4 GiB or more).
|
||||
let (direct, wanted): (Vec<_>, Vec<_>) = wanted.into_iter().partition(|c| {
|
||||
file_data.as_contiguous().is_none()
|
||||
&& pipeline.is_none_or(|pl| all_filters_skipped(pl, c.filter_mask))
|
||||
});
|
||||
for chunk in direct {
|
||||
read_unfiltered_overlap(
|
||||
file_data,
|
||||
chunk,
|
||||
chunk_bytes,
|
||||
&chunk_shape,
|
||||
elem_size,
|
||||
&mut boxed,
|
||||
&box_start,
|
||||
&box_extent,
|
||||
)?;
|
||||
}
|
||||
// Their stored bytes, batch by batch when the file is not in
|
||||
// memory; each batch's chunks are decoded into this thread's
|
||||
// reusable buffers before the next batch is fetched.
|
||||
|
||||
@@ -173,7 +173,7 @@ fn crafted() -> (Vec<u8>, Chunked, Vec<ChunkInfo>) {
|
||||
assert!(
|
||||
chunks
|
||||
.iter()
|
||||
.all(|c| c.chunk_size == HUGE && c.address == blob)
|
||||
.all(|c| c.chunk_size == u64::from(HUGE) && c.address == blob)
|
||||
);
|
||||
(bytes, ds, chunks)
|
||||
}
|
||||
@@ -313,7 +313,7 @@ fn large_reads_are_fetched_in_batches() {
|
||||
let data: Vec<u8> = (0..2 * CHUNK).map(|i| (i % 251) as u8).collect();
|
||||
let chunks: Vec<ChunkInfo> = (0..40u64)
|
||||
.map(|i| ChunkInfo {
|
||||
chunk_size: CHUNK as u32,
|
||||
chunk_size: CHUNK as u64,
|
||||
filter_mask: 0,
|
||||
offsets: vec![i * CHUNK as u64],
|
||||
address: (i % 2) * CHUNK as u64,
|
||||
|
||||
@@ -1528,11 +1528,11 @@ impl<'f> Dataset<'f> {
|
||||
let ds = self.dataspace()?;
|
||||
let dl = self.data_layout()?;
|
||||
let pipeline = self.filter_pipeline()?;
|
||||
// The selection reader knows nothing about fill values. When they
|
||||
// matter — no storage at all, or a non-zero fill on a chunked (possibly
|
||||
// sparse) dataset — select from a fill-aware full read instead. (The
|
||||
// selection reader currently decodes the full dataset too, so this
|
||||
// costs nothing extra.)
|
||||
// Fill values matter when there is no storage at all, or a non-zero
|
||||
// fill on a chunked (possibly sparse) dataset: a chunked dataset's
|
||||
// selection is then read over a box of fill values; anything else
|
||||
// (and a selection whose box is most of the dataset) is selected
|
||||
// from a fill-aware full read.
|
||||
let fill = clawhdf5_format::fill_value::dataset_fill_value_from_storage(
|
||||
&self.file.data,
|
||||
&self.header.messages,
|
||||
@@ -1545,6 +1545,29 @@ impl<'f> Dataset<'f> {
|
||||
&& !clawhdf5_format::fill_value::is_default(fill.as_deref()));
|
||||
if fill_matters {
|
||||
clawhdf5_format::partial_read::validate(selection, &ds.dimensions)?;
|
||||
// A chunked dataset: read the chunks the selection touches over
|
||||
// a box of fill values, not the whole dataset (which may be many
|
||||
// GiB when its chunks are large).
|
||||
if matches!(dl, DataLayout::Chunked { .. }) {
|
||||
clawhdf5_format::chunked_read::check_chunk_element_size(
|
||||
&dl,
|
||||
&dt,
|
||||
self.file.offset_size(),
|
||||
)?;
|
||||
if let Some(selected) = clawhdf5_format::partial_read::read_selection_filled_in(
|
||||
&self.file.data,
|
||||
&dl,
|
||||
&ds,
|
||||
dt.type_size() as usize,
|
||||
pipeline.as_ref(),
|
||||
self.file.offset_size(),
|
||||
self.file.length_size(),
|
||||
selection,
|
||||
Some(fill.as_deref().unwrap_or(&[])),
|
||||
)? {
|
||||
return Ok(selected);
|
||||
}
|
||||
}
|
||||
let full = self.read_raw()?;
|
||||
return Ok(data_read::extract_selection_from_buffer(
|
||||
&full,
|
||||
|
||||
+109
@@ -0,0 +1,109 @@
|
||||
"""Write HDF5 files whose chunks are 4 GiB or more, one dataset per chunk index.
|
||||
|
||||
python gen_huge_chunks.py filtered OUT.h5 # huge_chunks_filtered.h5
|
||||
python gen_huge_chunks.py unfiltered OUT.h5 # generated at test time
|
||||
|
||||
libhdf5 2.x writes a chunk of more than 0xFFFFFFFF bytes with layout message
|
||||
version 5 (`H5D__chunk_construct`: "chunk size > 4GB requires
|
||||
H5F_LIBVER_V200"), so never with a version-1 B-tree; a filtered chunk index
|
||||
element of a version-5 layout stores the chunk's size in "size of lengths"
|
||||
bytes (8). Every dataset is `<f8` with chunks of N = 2**29 + 1 elements
|
||||
(4 GiB + 8 bytes):
|
||||
|
||||
single shape (N,), chunks (N,) Single Chunk
|
||||
implicit shape (2N,), early allocation (unfiltered) Implicit
|
||||
farray shape (N+10,) Fixed Array
|
||||
earray shape (N+10,), maxshape (None,) Extensible Array
|
||||
btree2 shape (2, N+10), chunks (1, N),
|
||||
maxshape (None, None) v2 B-tree
|
||||
|
||||
Only a few elements are written: `d[0:10] = 0..9` and, where there is a
|
||||
second chunk along the axis, `d[N:N+10] = 100..109` (btree2: row 0 as that,
|
||||
row 1 `200..209` and `300..309`; implicit: `d[N:N+10]` and
|
||||
`d[2N-10:2N] = 500..509`). Everything else reads as the fill value.
|
||||
|
||||
`filtered` (the committed fixture, no `implicit`: that index is never
|
||||
filtered): deflate level 9 twice in the pipeline, fill value -1.0. One
|
||||
deflate leaves a 4 GiB chunk of a repeated 8-byte pattern at about 6 MiB;
|
||||
the second pass takes that to about 14 KiB, so the file is small while each
|
||||
chunk still inflates to 4 GiB + 8 bytes.
|
||||
|
||||
`unfiltered`: fill time "never" and the default fill value, so libhdf5
|
||||
writes only the elements written (an unfiltered chunk larger than the chunk
|
||||
cache is written in place) and the file is sparse: tens of GiB long but a few
|
||||
blocks on disk. Unwritten elements read as whatever the file holds there:
|
||||
zeros. Do not put it on tmpfs, which is memory.
|
||||
|
||||
libhdf5 holds a whole filtered chunk in memory while it writes it, so the
|
||||
`filtered` run needs about 4.5 GiB of memory (and about 3 minutes).
|
||||
|
||||
Generated 2026-09-28 with h5py 3.16.0 (libhdf5 2.0.0).
|
||||
"""
|
||||
import sys
|
||||
|
||||
import h5py
|
||||
import numpy as np
|
||||
|
||||
N = 2**29 + 1
|
||||
UNLIM = h5py.h5s.UNLIMITED
|
||||
|
||||
|
||||
def make(fid, name, index, filtered):
|
||||
dcpl = h5py.h5p.create(h5py.h5p.DATASET_CREATE)
|
||||
if index == "single":
|
||||
shape, maxshape, chunk = (N,), (N,), (N,)
|
||||
elif index == "implicit":
|
||||
shape, maxshape, chunk = (2 * N,), (2 * N,), (N,)
|
||||
dcpl.set_alloc_time(h5py.h5d.ALLOC_TIME_EARLY)
|
||||
elif index == "farray":
|
||||
shape, maxshape, chunk = (N + 10,), (N + 10,), (N,)
|
||||
elif index == "earray":
|
||||
shape, maxshape, chunk = (N + 10,), (UNLIM,), (N,)
|
||||
elif index == "btree2":
|
||||
shape, maxshape, chunk = (2, N + 10), (UNLIM, UNLIM), (1, N)
|
||||
dcpl.set_chunk(chunk)
|
||||
if filtered:
|
||||
dcpl.set_deflate(9)
|
||||
dcpl.set_deflate(9)
|
||||
dcpl.set_fill_value(np.array(-1.0, dtype="<f8"))
|
||||
else:
|
||||
dcpl.set_fill_time(h5py.h5d.FILL_TIME_NEVER)
|
||||
if index == "single":
|
||||
# libhdf5 2.0.0 drops a write to an unallocated unfiltered
|
||||
# Single Chunk this large when the fill time is "never" (the
|
||||
# dataset stays unallocated and reads zeros); allocated at
|
||||
# creation, the write lands in place.
|
||||
dcpl.set_alloc_time(h5py.h5d.ALLOC_TIME_EARLY)
|
||||
space = h5py.h5s.create_simple(shape, maxshape)
|
||||
dsid = h5py.h5d.create(fid, name.encode(), h5py.h5t.IEEE_F64LE, space, dcpl=dcpl)
|
||||
ds = h5py.Dataset(dsid)
|
||||
if len(shape) == 2:
|
||||
ds[0, 0:10] = np.arange(10.0)
|
||||
ds[0, N : N + 10] = np.arange(100.0, 110.0)
|
||||
ds[1, 0:10] = np.arange(200.0, 210.0)
|
||||
ds[1, N : N + 10] = np.arange(300.0, 310.0)
|
||||
else:
|
||||
ds[0:10] = np.arange(10.0)
|
||||
if shape[0] > N:
|
||||
ds[N : N + 10] = np.arange(100.0, 110.0)
|
||||
if index == "implicit":
|
||||
ds[2 * N - 10 : 2 * N] = np.arange(500.0, 510.0)
|
||||
dsid.close()
|
||||
|
||||
|
||||
def main():
|
||||
mode, out = sys.argv[1], sys.argv[2]
|
||||
filtered = {"filtered": True, "unfiltered": False}[mode]
|
||||
fapl = h5py.h5p.create(h5py.h5p.FILE_ACCESS)
|
||||
fapl.set_libver_bounds(h5py.h5f.LIBVER_EARLIEST, h5py.h5f.LIBVER_V200)
|
||||
fid = h5py.h5f.create(out.encode(), h5py.h5f.ACC_TRUNC, fapl=fapl)
|
||||
indexes = ["single", "farray", "earray", "btree2"]
|
||||
if not filtered:
|
||||
indexes.insert(1, "implicit")
|
||||
for index in indexes:
|
||||
make(fid, index, index, filtered)
|
||||
fid.close()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
Binary file not shown.
@@ -0,0 +1,296 @@
|
||||
//! Chunks of 4 GiB or more, which HDF5 2.0 writes (layout message version 5,
|
||||
//! `H5F_LIBVER_V200`), in every chunk index libhdf5 uses for them.
|
||||
//!
|
||||
//! `fixtures/huge_chunks_filtered.h5` (written by libhdf5 2.0.0 through h5py
|
||||
//! 3.16, `fixtures/gen_huge_chunks.py filtered`) holds one dataset per
|
||||
//! filtered index — Single Chunk, Fixed Array, Extensible Array, v2 B-tree —
|
||||
//! whose chunks are 2^29 + 1 `f64` (4 GiB + 8 bytes), stored deflated twice
|
||||
//! so each takes about 20 KiB. Listing their chunks is cheap and always runs;
|
||||
//! decoding one inflates 4 GiB, so those tests are opt-in:
|
||||
//! `CLAWHDF5_HUGE_CHUNKS=1` (run them one at a time, `--test-threads=1`:
|
||||
//! each needs about 4.5 GiB of memory). The same variable enables the
|
||||
//! unfiltered tests, which have h5py write a sparse file (about 44 GiB long,
|
||||
//! a few blocks on disk) under `tests/scratch/`, and the writer tests.
|
||||
//!
|
||||
//! libhdf5 never writes such a chunk with a version-1 B-tree (a chunk of more
|
||||
//! than 0xFFFFFFFF bytes forces layout version 5 and with it the newer
|
||||
//! indexes) and refuses to open one; `chunked_read`'s unit tests cover that.
|
||||
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::process::Command;
|
||||
|
||||
use clawhdf5::File;
|
||||
use clawhdf5_format::chunked_read::list_chunks;
|
||||
use clawhdf5_format::data_layout::DataLayout;
|
||||
use clawhdf5_format::dataspace::Dataspace;
|
||||
use clawhdf5_format::message_type::MessageType;
|
||||
use clawhdf5_format::object_header::ObjectHeader;
|
||||
use clawhdf5_format::selection::Selection;
|
||||
use clawhdf5_format::superblock::Superblock;
|
||||
|
||||
/// Elements per chunk along the chunked axis: 4 GiB + 8 bytes of `f64`.
|
||||
const N: u64 = (1 << 29) + 1;
|
||||
|
||||
fn fixture() -> PathBuf {
|
||||
Path::new(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures/huge_chunks_filtered.h5")
|
||||
}
|
||||
|
||||
fn heavy() -> bool {
|
||||
if std::env::var("CLAWHDF5_HUGE_CHUNKS").is_ok_and(|v| v == "1") {
|
||||
return true;
|
||||
}
|
||||
eprintln!("SKIP: set CLAWHDF5_HUGE_CHUNKS=1 to decode 4 GiB chunks");
|
||||
false
|
||||
}
|
||||
|
||||
fn python() -> String {
|
||||
std::env::var("CLAWHDF5_PYTHON").unwrap_or_else(|_| "python3".to_string())
|
||||
}
|
||||
|
||||
fn python_available() -> bool {
|
||||
Command::new(python())
|
||||
.args(["-c", "import h5py"])
|
||||
.output()
|
||||
.is_ok_and(|o| o.status.success())
|
||||
}
|
||||
|
||||
fn interop_required() -> bool {
|
||||
std::env::var("CLAWHDF5_REQUIRE_INTEROP").is_ok_and(|v| v == "1")
|
||||
}
|
||||
|
||||
/// The raw layout message version, the parsed layout, the dataspace and the
|
||||
/// chunks of `name` in the in-memory file `data`.
|
||||
fn layout_of(
|
||||
data: &[u8],
|
||||
name: &str,
|
||||
) -> (
|
||||
u8,
|
||||
DataLayout,
|
||||
Dataspace,
|
||||
Vec<clawhdf5_format::chunked_read::ChunkInfo>,
|
||||
) {
|
||||
let sb = Superblock::parse(data, 0).unwrap();
|
||||
let addr = clawhdf5_format::group_v2::resolve_path_any(data, &sb, name).unwrap();
|
||||
let hdr = ObjectHeader::parse(data, addr as usize, sb.offset_size, sb.length_size).unwrap();
|
||||
let msg = |t| {
|
||||
hdr.messages
|
||||
.iter()
|
||||
.find(|m| m.msg_type == t)
|
||||
.unwrap_or_else(|| panic!("{name}: no {t:?} message"))
|
||||
};
|
||||
let lm = msg(MessageType::DataLayout);
|
||||
let layout = DataLayout::parse(&lm.data, sb.offset_size, sb.length_size).unwrap();
|
||||
let space = Dataspace::parse(&msg(MessageType::Dataspace).data, sb.length_size).unwrap();
|
||||
let (chunks, _) =
|
||||
list_chunks(data, &layout, &space, 8, sb.offset_size, sb.length_size).unwrap();
|
||||
(lm.data[0], layout, space, chunks)
|
||||
}
|
||||
|
||||
/// The fixture's datasets: name, chunk index type, shape, and the scaled
|
||||
/// origins of the chunks libhdf5 wrote.
|
||||
const FILTERED: &[(&str, u8, &[u64], &[&[u64]])] = &[
|
||||
("single", 1, &[N], &[&[0]]),
|
||||
("farray", 3, &[N + 10], &[&[0], &[N]]),
|
||||
("earray", 4, &[N + 10], &[&[0], &[N]]),
|
||||
(
|
||||
"btree2",
|
||||
5,
|
||||
&[2, N + 10],
|
||||
&[&[0, 0], &[0, N], &[1, 0], &[1, N]],
|
||||
),
|
||||
];
|
||||
|
||||
/// Every index of the fixture: layout version 5, its chunks listed at the
|
||||
/// right origins with their stored (deflated) sizes. The index elements
|
||||
/// store a size in 8 bytes (libhdf5's `H5F_SIZEOF_SIZE` under layout
|
||||
/// version 5), where version 4 would use 6 for a chunk this size.
|
||||
#[test]
|
||||
fn filtered_huge_chunk_indexes_list() {
|
||||
let data = std::fs::read(fixture()).unwrap();
|
||||
for &(name, index, shape, origins) in FILTERED {
|
||||
let (version, layout, space, mut chunks) = layout_of(&data, name);
|
||||
assert_eq!(version, 5, "{name}: layout message version");
|
||||
let DataLayout::Chunked {
|
||||
chunk_dimensions,
|
||||
chunk_index_type,
|
||||
..
|
||||
} = &layout
|
||||
else {
|
||||
panic!("{name}: not chunked: {layout:?}");
|
||||
};
|
||||
assert_eq!(*chunk_index_type, Some(index), "{name}");
|
||||
assert_eq!(chunk_dimensions.last(), Some(&8), "{name}");
|
||||
assert_eq!(space.dimensions, shape, "{name}");
|
||||
chunks.sort_by(|a, b| a.offsets.cmp(&b.offsets));
|
||||
let got: Vec<&[u64]> = chunks.iter().map(|c| &c.offsets[..shape.len()]).collect();
|
||||
assert_eq!(got, origins, "{name}");
|
||||
for c in &chunks {
|
||||
assert!(
|
||||
(10_000..40_000).contains(&c.chunk_size),
|
||||
"{name}: stored size {} of chunk {:?}",
|
||||
c.chunk_size,
|
||||
c.offsets
|
||||
);
|
||||
assert_eq!(c.filter_mask, 0);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// `f64` values of `sel` in dataset `name` of `file`.
|
||||
fn sel(file: &File, name: &str, start: &[u64], count: &[u64]) -> Vec<f64> {
|
||||
let ds = file.dataset(name).unwrap();
|
||||
let s = Selection::Hyperslab {
|
||||
start: start.to_vec(),
|
||||
stride: vec![1; start.len()],
|
||||
count: count.to_vec(),
|
||||
block: vec![1; start.len()],
|
||||
};
|
||||
ds.read_f64_selection(&s).unwrap()
|
||||
}
|
||||
|
||||
fn f(range: std::ops::Range<i32>) -> Vec<f64> {
|
||||
range.map(f64::from).collect()
|
||||
}
|
||||
|
||||
/// Reads that touch a few elements of each 4 GiB chunk: written values, the
|
||||
/// fill value (-1) next to them, and the edge of the dataset.
|
||||
#[test]
|
||||
fn filtered_huge_chunks_read() {
|
||||
if !heavy() {
|
||||
return;
|
||||
}
|
||||
let file = File::open(fixture()).unwrap();
|
||||
let mut first = f(0..10);
|
||||
first.extend([-1.0; 2]);
|
||||
for name in ["single", "farray", "earray"] {
|
||||
assert_eq!(sel(&file, name, &[0], &[12]), first, "{name}");
|
||||
}
|
||||
assert_eq!(sel(&file, "single", &[N - 2], &[2]), [-1.0; 2]);
|
||||
let mut edge = vec![-1.0; 2];
|
||||
edge.extend(f(100..110));
|
||||
for name in ["farray", "earray"] {
|
||||
assert_eq!(sel(&file, name, &[N - 2], &[12]), edge, "{name}");
|
||||
}
|
||||
assert_eq!(sel(&file, "btree2", &[0, 0], &[1, 10]), f(0..10));
|
||||
assert_eq!(sel(&file, "btree2", &[1, N], &[1, 10]), f(300..310));
|
||||
assert_eq!(
|
||||
sel(&file, "btree2", &[0, N - 1], &[2, 2]),
|
||||
[-1.0, 100.0, -1.0, 300.0]
|
||||
);
|
||||
}
|
||||
|
||||
/// A directory under `tests/scratch/` (on disk: the sparse files must not
|
||||
/// land on a tmpfs `/tmp`), removed when dropped.
|
||||
fn scratch() -> tempfile::TempDir {
|
||||
let root = Path::new(env!("CARGO_MANIFEST_DIR")).join("tests/scratch");
|
||||
std::fs::create_dir_all(&root).unwrap();
|
||||
tempfile::tempdir_in(root).unwrap()
|
||||
}
|
||||
|
||||
/// Have h5py (libhdf5 2.x) write `fixtures/gen_huge_chunks.py`'s file for
|
||||
/// `mode` into `dir`; `None` when there is no h5py (a failure under
|
||||
/// `CLAWHDF5_REQUIRE_INTEROP=1`).
|
||||
fn generate(dir: &Path, mode: &str) -> Option<PathBuf> {
|
||||
if !python_available() {
|
||||
assert!(
|
||||
!interop_required(),
|
||||
"CLAWHDF5_REQUIRE_INTEROP=1 but python3 with h5py is not available"
|
||||
);
|
||||
eprintln!("SKIP: python3 with h5py not available");
|
||||
return None;
|
||||
}
|
||||
let script = Path::new(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures/gen_huge_chunks.py");
|
||||
let path = dir.join(format!("{mode}.h5"));
|
||||
let out = Command::new(python())
|
||||
.arg(&script)
|
||||
.arg(mode)
|
||||
.arg(&path)
|
||||
.output()
|
||||
.unwrap();
|
||||
assert!(
|
||||
out.status.success(),
|
||||
"gen_huge_chunks.py {mode} failed:\n{}",
|
||||
String::from_utf8_lossy(&out.stderr)
|
||||
);
|
||||
Some(path)
|
||||
}
|
||||
|
||||
/// Positioned reads of a file (no mmap), counting the bytes read.
|
||||
struct Counting {
|
||||
file: std::fs::File,
|
||||
len: u64,
|
||||
read: std::sync::atomic::AtomicU64,
|
||||
}
|
||||
|
||||
impl clawhdf5_format::storage::Storage for Counting {
|
||||
fn read_at(
|
||||
&self,
|
||||
offset: u64,
|
||||
len: usize,
|
||||
) -> Result<std::borrow::Cow<'_, [u8]>, clawhdf5_format::error::FormatError> {
|
||||
use std::os::unix::fs::FileExt;
|
||||
let len = len.min(usize::try_from(self.len.saturating_sub(offset)).unwrap_or(usize::MAX));
|
||||
let mut buf = vec![0; len];
|
||||
self.file.read_exact_at(&mut buf, offset).unwrap();
|
||||
self.read
|
||||
.fetch_add(len as u64, std::sync::atomic::Ordering::Relaxed);
|
||||
Ok(buf.into())
|
||||
}
|
||||
|
||||
fn len(&self) -> u64 {
|
||||
self.len
|
||||
}
|
||||
}
|
||||
|
||||
fn check_unfiltered(file: &File) {
|
||||
let mut first = f(0..10);
|
||||
first.extend([0.0; 2]);
|
||||
for name in ["single", "implicit", "farray", "earray"] {
|
||||
assert_eq!(sel(file, name, &[0], &[12]), first, "{name}");
|
||||
}
|
||||
assert_eq!(sel(file, "single", &[N - 2], &[2]), [0.0; 2]);
|
||||
let mut edge = vec![0.0; 2];
|
||||
edge.extend(f(100..110));
|
||||
for name in ["implicit", "farray", "earray"] {
|
||||
assert_eq!(sel(file, name, &[N - 2], &[12]), edge, "{name}");
|
||||
}
|
||||
let mut last = vec![0.0; 2];
|
||||
last.extend(f(500..510));
|
||||
assert_eq!(sel(file, "implicit", &[2 * N - 12], &[12]), last);
|
||||
assert_eq!(sel(file, "btree2", &[0, 0], &[1, 10]), f(0..10));
|
||||
assert_eq!(sel(file, "btree2", &[1, N], &[1, 10]), f(300..310));
|
||||
assert_eq!(
|
||||
sel(file, "btree2", &[0, N - 1], &[2, 2]),
|
||||
[0.0, 100.0, 0.0, 300.0]
|
||||
);
|
||||
}
|
||||
|
||||
/// Unfiltered chunks of 4 GiB + 8 bytes in every index libhdf5 gives them
|
||||
/// (Single Chunk, Implicit, Fixed Array, Extensible Array, v2 B-tree), in a
|
||||
/// sparse file h5py writes: read through a memory map, and through
|
||||
/// positioned reads, where a selection reads only the rows it needs.
|
||||
#[test]
|
||||
fn unfiltered_huge_chunks_read() {
|
||||
if !heavy() {
|
||||
return;
|
||||
}
|
||||
let dir = scratch();
|
||||
let Some(path) = generate(dir.path(), "unfiltered") else {
|
||||
return;
|
||||
};
|
||||
check_unfiltered(&File::open(&path).unwrap());
|
||||
|
||||
let file = std::fs::File::open(&path).unwrap();
|
||||
let len = file.metadata().unwrap().len();
|
||||
assert!(len > 40 << 30, "{len}");
|
||||
let storage = std::sync::Arc::new(Counting {
|
||||
file,
|
||||
len,
|
||||
read: 0.into(),
|
||||
});
|
||||
let positioned = File::open_storage(storage.clone()).unwrap();
|
||||
check_unfiltered(&positioned);
|
||||
// Every read above together: metadata and a few rows, not 4 GiB chunks.
|
||||
let read = storage.read.load(std::sync::atomic::Ordering::Relaxed);
|
||||
assert!(read < 1 << 20, "{read} bytes read");
|
||||
}
|
||||
@@ -283,7 +283,7 @@ mod parallel_tests {
|
||||
for (i, chunk) in compressed_chunks.iter().enumerate() {
|
||||
file_data[offset..offset + chunk.len()].copy_from_slice(chunk);
|
||||
chunk_infos.push(ChunkInfo {
|
||||
chunk_size: chunk.len() as u32,
|
||||
chunk_size: chunk.len() as u64,
|
||||
filter_mask: 0,
|
||||
offsets: vec![(i * chunk_elems) as u64, 0],
|
||||
address: offset as u64,
|
||||
|
||||
Reference in New Issue
Block a user