Streaming minibatches
Mix.install([
{:ex_zarr, "~> 1.3"},
# {:ex_zarr, path: Path.join(__DIR__, "../../..")},
{:nx, "~> 0.7"}
])
Setup
This notebook builds a tiny feature matrix and a label vector, then reads them three ways.
- One tensor per stored chunk. The grid is whatever you set in
chunks. - One tensor per minibatch. You choose how many samples (axis 0) each step sees. A batch can cross a chunk boundary.
- Paired feature and label batches, with a shuffle that is repeatable when you pass
:seed.
A 10×4 feature matrix and 10 labels
features has shape {10, 4} and chunks {4, 4}. Along the sample axis that is three chunks: rows 0–3, rows 4–7, and rows 8–9. The last chunk is shorter than 4 rows.
labels is a length-10 vector with chunks of 4, so it has the same sample boundaries. Values are just the index, stored as little-endian float32, which is the layout Nx.from_binary/2 expects for {:f, 32}.
alias ExZarr.Array
{:ok, features} =
Array.create(shape: {10, 4}, chunks: {4, 4}, dtype: :float32, storage: :memory)
{:ok, labels} =
Array.create(shape: {10}, chunks: {4}, dtype: :float32, storage: :memory)
x = for i <- 0..39, into: <<>>, do: <<i * 1.0::float-little-32>>
y = for i <- 0..9, into: <<>>, do: <<i * 1.0::float-little-32>>
:ok = Array.set_slice(features, x, start: {0, 0}, stop: {10, 4})
:ok = Array.set_slice(labels, y, start: {0}, stop: {10})
{features.metadata.shape, features.chunks, labels.metadata.shape, labels.chunks}
One tensor per stored chunk
stream_chunk_tensors/2 follows the chunk grid. It does not invent a batch size.
:concurrency is forwarded to Array.stream_chunks/2 (here 2 readers). :ordered keeps chunk index order, so the shapes come out as {4, 4}, {4, 4}, then the edge {2, 4}. The edge chunk is stored padded to {4, 4}; the helper crops it to the two rows that belong to the array.
Use this path for reductions and maps that can run per chunk. The tensor shape changes on the last chunk, so a training step that requires a static batch size should use DataLoader instead.
features
|> ExZarr.Nx.stream_chunk_tensors(concurrency: 2, ordered: true)
|> Enum.map(fn {:ok, tensor} -> Nx.shape(tensor) end)
Minibatches of 3 samples
batch_stream/3 slices axis 0 with get_slice/2. Each item is {:ok, tensor} of shape {3, 4}: 3 samples, all 4 features.
drop_remainder: true drops the final incomplete group. 10 samples ÷ 3 leaves 1 sample, so you get three batches and the last sample is omitted. The second batch starts at row 3, which is inside the first chunk, and ends at row 6, which is inside the second chunk. That single batch reads across the chunk boundary. A chunk stream cannot do that.
features
|> ExZarr.Nx.DataLoader.batch_stream(3, drop_remainder: true)
|> Enum.map(fn {:ok, batch} -> Nx.shape(batch) end)
Shuffled feature/label pairs
paired_shuffled_batch_stream/4 draws the same sample indexes for features and labels, so row i of a feature batch is still label i.
The shuffle builds the full index list 0..9, permutes it, then loads each batch. :seed makes that permutation repeatable. :shuffle_buffer_size is accepted and ignored. There is no partial buffer.
Batch size 4 with 10 samples yields shapes {4, 4} / {4}, then {4, 4} / {4}, then {2, 4} / {2}. The short last pair is the remainder. Pass drop_remainder: true when the training step needs every batch to have the same size.
ExZarr.Nx.DataLoader.paired_shuffled_batch_stream(features, labels, 4, seed: 1)
|> Enum.map(fn {:ok, {xb, yb}} -> {Nx.shape(xb), Nx.shape(yb)} end)