From 48e526860e01f6a81a6f9add73b5851d5efdf73d Mon Sep 17 00:00:00 2001 From: tpietzsch Date: Mon, 2 Feb 2026 21:16:29 +0100 Subject: [PATCH 01/10] Add SliceTrackingLazyRead This is in preparation for implementing prefetch() for cloud storage backends... --- .../prefetch/SliceTrackingLazyRead.java | 114 +++++++++++++++ .../n5/readdata/prefetch/Slices.java | 116 +++++++++++++++ .../n5/readdata/prefetch/SlicesTest.java | 134 ++++++++++++++++++ 3 files changed, 364 insertions(+) create mode 100644 src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyRead.java create mode 100644 src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/Slices.java create mode 100644 src/test/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SlicesTest.java diff --git a/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyRead.java b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyRead.java new file mode 100644 index 000000000..3ffd852ea --- /dev/null +++ b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyRead.java @@ -0,0 +1,114 @@ +package org.janelia.saalfeldlab.n5.readdata.prefetch; + +import java.io.IOException; +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; +import org.janelia.saalfeldlab.n5.N5Exception.N5IOException; +import org.janelia.saalfeldlab.n5.readdata.LazyRead; +import org.janelia.saalfeldlab.n5.readdata.Range; +import org.janelia.saalfeldlab.n5.readdata.ReadData; + +public class SliceTrackingLazyRead implements LazyRead { + + private static class Slice implements Range { + + // Offset and length in the delegate + private final long offset; + private final long length; + + // Data of this slice + private final ReadData data; + + Slice(final long offset, final long length, final ReadData data) { + this.offset = offset; + this.length = length; + this.data = data; + } + + @Override + public long offset() { + return offset; + } + + @Override + public long length() { + return length; + } + + @Override + public String toString() { + return "{" + offset + ", " + length + '}'; + } + } + + private final List slices = new ArrayList<>(); + + /** + * The {@code LazyRead} providing our data. + */ + private final LazyRead delegate; + + public SliceTrackingLazyRead(final LazyRead delegate) { + this.delegate = delegate; + } + + @Override + public void close() throws IOException { + delegate.close(); + } + + @Override + public ReadData materialize(final long offset, final long length) throws N5IOException { + final Slice containing = Slices.findContainingSlice(slices, offset, length); + if (containing != null) { + return containing.data.slice(offset - containing.offset, length); + } else { + final ReadData data = delegate.materialize(offset, length); + Slices.addSlice(slices, new Slice(offset, length, data)); + return data; + } + } + + @Override + public long size() throws N5IOException { + return delegate.size(); + } + + /** + * Indicates that the given slices will be subsequently read. + * {@code LazyRead} implementations (optionally) may take steps to prepare + * for these subsequent slices. + *

+ * Minimal implementation: Find offset and length covering all ranges that + * are not yet fully covered by existing slices. Then materialize the slice + * covering that range. + * + * @param ranges + * slice ranges to prefetch + * + * @throws N5IOException + * if any I/O error occurs + */ + @Override + public void prefetch(final Collection ranges) throws N5IOException { + + long fromIndex = Long.MAX_VALUE; + long toIndex = Long.MIN_VALUE; + for (final Range slice : ranges) { + if (!isCovered(slice)) { + fromIndex = Math.min(fromIndex, slice.offset()); + toIndex = Math.max(toIndex, slice.end()); + } + } + + if (fromIndex < toIndex) { + materialize(fromIndex, toIndex - fromIndex); + } + } + + private boolean isCovered(final Range slice) { + + return Slices.findContainingSlice(slices, slice) != null; + } +} diff --git a/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/Slices.java b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/Slices.java new file mode 100644 index 000000000..d3d0dfb11 --- /dev/null +++ b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/Slices.java @@ -0,0 +1,116 @@ +package org.janelia.saalfeldlab.n5.readdata.prefetch; + +import java.util.Collections; +import java.util.Comparator; +import java.util.List; +import org.janelia.saalfeldlab.n5.readdata.Range; + +class Slices { + + private Slices() { + // utility class. should not be instantiated. + } + + /** + * In an ordered list of {@code slices}, find a slice that completely contains the given range. + *

+ * Pre-conditions: + *

    + *
  1. Slices are ordered by offset.
  2. + *
  3. If two slices overlap, no slice is fully contained within the other. + * (Therefore, if {@code a.offset < b.offset} then {@code a.end < b.end}.)
  4. + *
+ * + * @param slices + * ordered list of slices + * @param offset + * start of the range to cover + * @param length + * length of the range to cover + * + * @return a slice that completely contains the requested range, or {@code null} if no such slice exists + */ + static T findContainingSlice(final List slices, final long offset, final long length) { + + // Find the slice with the largest slice.offset <= offset. + final int i = Collections.binarySearch(slices, Range.at(offset, 0), Comparator.comparingLong(Range::offset)); + + // Largest index of a slice with slice.offset <= offset. + final int index = i < 0 ? -i - 2 : i; + if (index < 0) { + // We find no overlapping slice, because + // slices[0].offset is already too large. + return null; + } + + final T slice = slices.get(index); + if (slice.end() < offset + length) { + return null; + } + + return slice; + } + + /** + * In an ordered list of {@code slices}, find a slice that completely contains the given range. + *

+ * Pre-conditions: + *

    + *
  1. Slices are ordered by offset.
  2. + *
  3. If two slices overlap, no slice is fully contained within the other. + * (Therefore, if {@code a.offset < b.offset} then {@code a.end < b.end}.)
  4. + *
+ * + * @param slices + * ordered list of slices + * @param range + * range to cover + * + * @return a slice that completely contains the requested range, or {@code null} if no such slice exists + */ + static T findContainingSlice(final List slices, final Range range) { + return findContainingSlice(slices, range.offset(), range.length()); + } + + /** + * Add a new {@code slice} to the {@code slice} list. + *

+ * Note, that the new {@code slice} is expected to not be fully contained in + * an existing slice! + *

+ * Pre/post-conditions: + *

    + *
  1. Slices are ordered by offset.
  2. + *
  3. If two slices overlap, no slice is fully contained within the other. + * (Therefore, if {@code a.offset < b.offset} then {@code a.end < b.end}.)
  4. + *
+ *

+ * The new {@code slice} will be inserted into the list at the correct position + * (such that {@code slices} remains ordered by slice offset), and all existing + * slices that are fully contained in the new {@code slice} will be removed. + * + * @param slices + * ordered list of slices + * @param slice + * slice to be inserted + */ + static void addSlice(final List slices, final T slice) { + + final int i = Collections.binarySearch(slices, slice, Comparator.comparingLong(Range::offset)); + final int from = i < 0 ? -i - 1 : i; + + int to = from; + while (to < slices.size() && slices.get(to).end() <= slice.end()) { + ++to; + } + + if (from == to) { + // empty range: just insert + slices.add(from, slice); + } else { + // overwrite the first element in range, remove the rest + slices.set(from, slice); + slices.subList(from + 1, to).clear(); + } + } +} diff --git a/src/test/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SlicesTest.java b/src/test/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SlicesTest.java new file mode 100644 index 000000000..0ce2b7442 --- /dev/null +++ b/src/test/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SlicesTest.java @@ -0,0 +1,134 @@ +package org.janelia.saalfeldlab.n5.readdata.prefetch; + +import java.util.ArrayList; +import java.util.List; +import org.janelia.saalfeldlab.n5.readdata.Range; +import org.junit.Test; + +import static org.junit.Assert.assertEquals; + + +public class SlicesTest { + + private List createSlices(final long[] offsets, final long[] lengths) { + final List slices = new ArrayList<>(); + for (int i = 0; i < offsets.length; ++i) { + slices.add(Range.at(offsets[i], lengths[i])); + } + return slices; + } + + @Test + public void testFindContaining() { + + // 0 1 2 3 4 5 6 7 8 9 A B C D E F + // (2,6) [-----------] + // (6,4) [---------] + // (8,6) [-----------] + + final List slices = createSlices( + new long[] {2, 6, 8}, + new long[] {6, 4, 6}); + Range slice; + + // 0 1 2 3 4 5 6 7 8 9 A B C D E F + // (1,1) [-] + slice = Slices.findContainingSlice(slices, 1, 1); + assertEquals(null, slice); + + // 0 1 2 3 4 5 6 7 8 9 A B C D E F + // (2,1) [-] + slice = Slices.findContainingSlice(slices, 2, 1); + assertEquals(2, slice.offset()); + + // 0 1 2 3 4 5 6 7 8 9 A B C D E F + // (2,6) [-----------] + slice = Slices.findContainingSlice(slices, 2, 6); + assertEquals(2, slice.offset()); + + // 0 1 2 3 4 5 6 7 8 9 A B C D E F + // (2,7) [-------------] + slice = Slices.findContainingSlice(slices, 2, 7); + assertEquals(null, slice); + + // 0 1 2 3 4 5 6 7 8 9 A B C D E F + // (6,4) [-------] + slice = Slices.findContainingSlice(slices, 6, 4); + assertEquals(6, slice.offset()); + + // 0 1 2 3 4 5 6 7 8 9 A B C D E F + // (8,2) [---] + slice = Slices.findContainingSlice(slices, 8, 2); + assertEquals(8, slice.offset()); + + // 0 1 2 3 4 5 6 7 8 9 A B C D E F + // (12,2) [---] + slice = Slices.findContainingSlice(slices, 12, 2); + assertEquals(8, slice.offset()); + + // 0 1 2 3 4 5 6 7 8 9 A B C D E F + // (12,3) [-----] + slice = Slices.findContainingSlice(slices, 12, 3); + assertEquals(null, slice); + + // 0 1 2 3 4 5 6 7 8 9 A B C D E F + // (14,1) [-] + slice = Slices.findContainingSlice(slices, 14, 1); + assertEquals(null, slice); + } + + + @Test + public void testAddSlice() { + + // 0 1 2 3 4 5 6 7 8 9 A B C D E F + // (2,6) [-----------] + // (6,4) [---------] + // (8,6) [-----------] + final List initial = createSlices( + new long[] {2, 6, 8}, + new long[] {6, 4, 6}); + List slices; + + + slices = new ArrayList<>(initial); + Slices.addSlice(slices, Range.at(0, 1)); + // 0 1 2 3 4 5 6 7 8 9 A B C D E F + // (0,1) [-] + // (2,6) [-----------] + // (6,4) [---------] + // (8,6) [-----------] + assertEquals(createSlices( + new long[] {0, 2, 6, 8}, + new long[] {1, 6, 4, 6}), slices); + + + slices = new ArrayList<>(initial); + Slices.addSlice(slices, Range.at(0, 16)); + // 0 1 2 3 4 5 6 7 8 9 A B C D E F + // (0,16)[-------------------------------] + assertEquals(createSlices( + new long[] {0}, + new long[] {16}), slices); + + + slices = new ArrayList<>(initial); + Slices.addSlice(slices, Range.at(2, 8)); + // 0 1 2 3 4 5 6 7 8 9 A B C D E F + // (2,8) [-----------------] + // (8,6) [-----------] + assertEquals(createSlices( + new long[] {2, 8}, + new long[] {8, 6}), slices); + + + slices = new ArrayList<>(initial); + Slices.addSlice(slices, Range.at(1, 10)); + // 0 1 2 3 4 5 6 7 8 9 A B C D E F + // (1,10) [---------------------] + // (8,6) [-----------] + assertEquals(createSlices( + new long[] {1, 8}, + new long[] {10, 6}), slices); + } +} From 7853ca697b588eea47cd10574138f4d29e14aa14 Mon Sep 17 00:00:00 2001 From: John Bogovic Date: Mon, 2 Feb 2026 17:00:34 -0500 Subject: [PATCH 02/10] feat: add Range.aggregate * and a test --- .../saalfeldlab/n5/readdata/Range.java | 50 +++++++++++++++- .../saalfeldlab/n5/readdata/RangeTests.java | 58 +++++++++++++++++++ 2 files changed, 107 insertions(+), 1 deletion(-) create mode 100644 src/test/java/org/janelia/saalfeldlab/n5/readdata/RangeTests.java diff --git a/src/main/java/org/janelia/saalfeldlab/n5/readdata/Range.java b/src/main/java/org/janelia/saalfeldlab/n5/readdata/Range.java index eff501ba2..a42c18ba4 100644 --- a/src/main/java/org/janelia/saalfeldlab/n5/readdata/Range.java +++ b/src/main/java/org/janelia/saalfeldlab/n5/readdata/Range.java @@ -28,6 +28,9 @@ */ package org.janelia.saalfeldlab.n5.readdata; +import java.util.ArrayList; +import java.util.Collection; +import java.util.Collections; import java.util.Comparator; /** @@ -61,7 +64,7 @@ default long end() { } static boolean equals(final Range r0, final Range r1) { - if (r0 == null && r1==null) { + if (r0 == null && r1 == null) { return true; } else if (r0 == null || r1 == null) { return false; @@ -119,4 +122,49 @@ public String toString() { return new DefaultRange(offset, length); } + + /** + * Returns a potentially new collection of {@link Range}s such that + * adjacent or overlapping Ranges are combined. + *

+ * If the input ranges are non-adjacent the input instance is returned, but if aggregation + * occurs, the result will be a new, sorted List. + * + * @param ranges + * a collection of Ranges + * @return + */ + static Collection aggregate(final Collection ranges) { + + if (ranges.size() == 0) + return ranges; + + ArrayList sortedRanges = new ArrayList<>(ranges); + Collections.sort(sortedRanges, Range.COMPARATOR); + + ArrayList result = new ArrayList<>(); + boolean wereMerges = false; + Range lo = null; + for (Range hi : sortedRanges) { + + if (lo == null) + lo = hi; + else if (lo.end() >= hi.offset()) { + // merge + final Range mergedLo = Range.at(lo.offset(), Math.max(lo.end(), hi.end()) - lo.offset()); + lo = mergedLo; + wereMerges = true; + } else { + result.add(lo); + lo = hi; + } + } + result.add(lo); + + if (!wereMerges) + return ranges; + + return result; + } + } diff --git a/src/test/java/org/janelia/saalfeldlab/n5/readdata/RangeTests.java b/src/test/java/org/janelia/saalfeldlab/n5/readdata/RangeTests.java new file mode 100644 index 000000000..3499220b9 --- /dev/null +++ b/src/test/java/org/janelia/saalfeldlab/n5/readdata/RangeTests.java @@ -0,0 +1,58 @@ +package org.janelia.saalfeldlab.n5.readdata; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertSame; + +import java.util.Collections; +import java.util.List; +import java.util.stream.Collectors; +import java.util.stream.Stream; + +import org.junit.Test; + +public class RangeTests { + + @Test + public void testAggregate() { + + List nonOverlapping = Stream.of(Range.at(0, 5), Range.at(12, 2)).collect(Collectors.toList()); + assertSame(nonOverlapping, Range.aggregate(nonOverlapping)); + + List nonOverlappingRev = Stream.of(Range.at(12, 2), Range.at(0, 5)).collect(Collectors.toList()); + assertSame(nonOverlappingRev, Range.aggregate(nonOverlappingRev)); + + /** + * 0 1 2 3 4 5 + * x x x x x - + * - x - - - - + */ + List containing = Stream.of(Range.at(0, 5), Range.at(1, 1)).collect(Collectors.toList()); + assertEquals(Collections.singletonList(Range.at(0, 5)), Range.aggregate(containing)); + + /** + * 0 1 2 3 4 5 6 + * x x x x x - - + * - - x x x x x + */ + List overlapping = Stream.of(Range.at(0, 5), Range.at(2, 5)).collect(Collectors.toList()); + assertEquals(Collections.singletonList(Range.at(0, 7)), Range.aggregate(overlapping)); + + /** + * 0 1 2 3 4 5 + * x x x x x - + * - - - - - x + */ + List adjacent = Stream.of(Range.at(0, 5), Range.at(5, 1)).collect(Collectors.toList()); + assertEquals(Collections.singletonList(Range.at(0, 6)), Range.aggregate(adjacent)); + + /** + * 0 1 2 3 4 5 + * - - - x x - + * x x - - - - + * - - - - - x + */ + List three = Stream.of(Range.at(3, 2), Range.at(0, 2), Range.at(5, 1)).collect(Collectors.toList()); + assertEquals(Stream.of(Range.at(0, 2), Range.at(3, 3)).collect(Collectors.toList()), Range.aggregate(three)); + } + +} From b2e871df95448237bf14ad49fb194f45bef6993b Mon Sep 17 00:00:00 2001 From: tpietzsch Date: Fri, 3 Apr 2026 21:40:56 +0200 Subject: [PATCH 03/10] simplify --- .../saalfeldlab/n5/readdata/Range.java | 25 ++++++++----------- 1 file changed, 10 insertions(+), 15 deletions(-) diff --git a/src/main/java/org/janelia/saalfeldlab/n5/readdata/Range.java b/src/main/java/org/janelia/saalfeldlab/n5/readdata/Range.java index a42c18ba4..06e753366 100644 --- a/src/main/java/org/janelia/saalfeldlab/n5/readdata/Range.java +++ b/src/main/java/org/janelia/saalfeldlab/n5/readdata/Range.java @@ -30,7 +30,6 @@ import java.util.ArrayList; import java.util.Collection; -import java.util.Collections; import java.util.Comparator; /** @@ -131,19 +130,19 @@ public String toString() { * occurs, the result will be a new, sorted List. * * @param ranges - * a collection of Ranges - * @return + * a collection of Ranges + * + * @return collection with adjacent or overlapping ranges merged */ static Collection aggregate(final Collection ranges) { if (ranges.size() == 0) return ranges; - ArrayList sortedRanges = new ArrayList<>(ranges); - Collections.sort(sortedRanges, Range.COMPARATOR); + final ArrayList sortedRanges = new ArrayList<>(ranges); + sortedRanges.sort(Range.COMPARATOR); - ArrayList result = new ArrayList<>(); - boolean wereMerges = false; + final ArrayList result = new ArrayList<>(); Range lo = null; for (Range hi : sortedRanges) { @@ -151,9 +150,7 @@ static Collection aggregate(final Collection r lo = hi; else if (lo.end() >= hi.offset()) { // merge - final Range mergedLo = Range.at(lo.offset(), Math.max(lo.end(), hi.end()) - lo.offset()); - lo = mergedLo; - wereMerges = true; + lo = Range.at(lo.offset(), Math.max(lo.end(), hi.end()) - lo.offset()); } else { result.add(lo); lo = hi; @@ -161,10 +158,8 @@ else if (lo.end() >= hi.offset()) { } result.add(lo); - if (!wereMerges) - return ranges; - - return result; - } + final boolean wereMerges = ranges.size() != result.size(); + return wereMerges ? result : ranges; + } } From 4d56228c093e53f5c9d2686031a50396e98bd859 Mon Sep 17 00:00:00 2001 From: John Bogovic Date: Mon, 2 Feb 2026 17:01:52 -0500 Subject: [PATCH 04/10] feat: make SliceTrackingLazyRead abstract, add two implementations * Default- and Aggregating- * add tests --- .../AggregatingSliceTrackingLazyRead.java | 35 +++ .../DefaultSliceTrackingLazyRead.java | 50 ++++ .../prefetch/SliceTrackingLazyRead.java | 6 +- .../prefetch/SliceTrackingLazyReadTests.java | 233 ++++++++++++++++++ 4 files changed, 321 insertions(+), 3 deletions(-) create mode 100644 src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/AggregatingSliceTrackingLazyRead.java create mode 100644 src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/DefaultSliceTrackingLazyRead.java create mode 100644 src/test/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyReadTests.java diff --git a/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/AggregatingSliceTrackingLazyRead.java b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/AggregatingSliceTrackingLazyRead.java new file mode 100644 index 000000000..fc4fb660d --- /dev/null +++ b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/AggregatingSliceTrackingLazyRead.java @@ -0,0 +1,35 @@ +package org.janelia.saalfeldlab.n5.readdata.prefetch; + +import java.util.Collection; + +import org.janelia.saalfeldlab.n5.N5Exception.N5IOException; +import org.janelia.saalfeldlab.n5.readdata.LazyRead; +import org.janelia.saalfeldlab.n5.readdata.Range; + +public class AggregatingSliceTrackingLazyRead extends SliceTrackingLazyRead { + + public AggregatingSliceTrackingLazyRead(final LazyRead delegate) { + super(delegate); + } + + /** + * Indicates that the given slices will be subsequently read. + *

+ * This implementation groups overlapping / adjacent {@link Range}s into single read requests. + * + * @param ranges + * slice ranges to prefetch + * + * @throws N5IOException + * if any I/O error occurs + */ + @Override + public void prefetch(final Collection ranges) throws N5IOException { + + final Collection aggregatedRanges = Range.aggregate(ranges); + for (final Range slice : aggregatedRanges) { + materialize(slice.offset(), slice.length()); + } + } + +} diff --git a/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/DefaultSliceTrackingLazyRead.java b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/DefaultSliceTrackingLazyRead.java new file mode 100644 index 000000000..1e607ddc1 --- /dev/null +++ b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/DefaultSliceTrackingLazyRead.java @@ -0,0 +1,50 @@ +package org.janelia.saalfeldlab.n5.readdata.prefetch; + +import java.util.Collection; +import org.janelia.saalfeldlab.n5.N5Exception.N5IOException; +import org.janelia.saalfeldlab.n5.readdata.LazyRead; +import org.janelia.saalfeldlab.n5.readdata.Range; + +public class DefaultSliceTrackingLazyRead extends SliceTrackingLazyRead { + + public DefaultSliceTrackingLazyRead(final LazyRead delegate) { + super(delegate); + } + + /** + * Indicates that the given slices will be subsequently read. + * {@code LazyRead} implementations (optionally) may take steps to prepare + * for these subsequent slices. + *

+ * Minimal implementation: Find offset and length covering all ranges that + * are not yet fully covered by existing slices. Then materialize the slice + * covering that range. + * + * @param ranges + * slice ranges to prefetch + * + * @throws N5IOException + * if any I/O error occurs + */ + @Override + public void prefetch(final Collection ranges) throws N5IOException { + + long fromIndex = Long.MAX_VALUE; + long toIndex = Long.MIN_VALUE; + for (final Range slice : ranges) { + if (!isCovered(slice)) { + fromIndex = Math.min(fromIndex, slice.offset()); + toIndex = Math.max(toIndex, slice.end()); + } + } + + if (fromIndex < toIndex) { + materialize(fromIndex, toIndex - fromIndex); + } + } + + private boolean isCovered(final Range slice) { + + return Slices.findContainingSlice(slices, slice) != null; + } +} diff --git a/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyRead.java b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyRead.java index 3ffd852ea..de20f2c39 100644 --- a/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyRead.java +++ b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyRead.java @@ -9,9 +9,9 @@ import org.janelia.saalfeldlab.n5.readdata.Range; import org.janelia.saalfeldlab.n5.readdata.ReadData; -public class SliceTrackingLazyRead implements LazyRead { +public abstract class SliceTrackingLazyRead implements LazyRead { - private static class Slice implements Range { + protected static class Slice implements Range { // Offset and length in the delegate private final long offset; @@ -42,7 +42,7 @@ public String toString() { } } - private final List slices = new ArrayList<>(); + protected final List slices = new ArrayList<>(); /** * The {@code LazyRead} providing our data. diff --git a/src/test/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyReadTests.java b/src/test/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyReadTests.java new file mode 100644 index 000000000..044e6c1a3 --- /dev/null +++ b/src/test/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyReadTests.java @@ -0,0 +1,233 @@ +package org.janelia.saalfeldlab.n5.readdata.prefetch; + +import static org.junit.Assert.assertEquals; + +import java.io.IOException; +import java.util.Arrays; +import java.util.List; + +import org.janelia.saalfeldlab.n5.N5Exception.N5IOException; +import org.janelia.saalfeldlab.n5.readdata.LazyRead; +import org.janelia.saalfeldlab.n5.readdata.Range; +import org.janelia.saalfeldlab.n5.readdata.ReadData; +import org.junit.Test; + +public class SliceTrackingLazyReadTests { + + @Test + public void testDefaultSliceTracking() throws N5IOException { + + /** + * 1. Create sample ReadData from byte[] + * 2. Create a DummyLazy Read + * 3. Make a DefaultSliceTrackingLazyRead + * 4. Create a list of ranges + * 5. Call prefetch + * 6. Ensure the correct number of materialize calls were made (always 1 for DefaultSliceTrackingLazyRead) + * 7. Verify the stored slices contain the correct range + */ + + // 1-2. Create a DummyLazyRead with 64 bytes + DummyLazyRead dummyLazyRead = createDummyLazyRead(64); + + // 3. Make a testable DefaultSliceTrackingLazyRead + TestableDefaultSliceTracker sliceTracking = new TestableDefaultSliceTracker(dummyLazyRead); + + // 4. Create a list of ranges (two non-overlapping ranges with a gap) + List ranges = Arrays.asList( + Range.at(10, 5), // offset 10, length 5 (bytes 10-14) + Range.at(50, 10) // offset 50, length 10 (bytes 50-59) + ); + + // 5. Call prefetch + sliceTracking.prefetch(ranges); + + // 6. Ensure exactly 1 materialize call was made + // DefaultSliceTrackingLazyRead creates a single large slice covering all ranges + assertEquals("DefaultSliceTrackingLazyRead should make exactly 1 materialize call", + 1, dummyLazyRead.getNumMaterializeCalls()); + + // 7. Verify the stored slice covers the entire range from offset 10 with length 50 + assertStoredSlices(sliceTracking, Arrays.asList( + Range.at(10, 50) // Single slice covering offset 10-59 + )); + + } + + @Test + public void testAggregatingSliceTracking() throws N5IOException { + + /** + * 1. Create sample ReadData from byte[] + * 2. Create a DummyLazyRead + * 3. Make an AggregatingSliceTrackingLazyRead + * 4. Create a list of ranges + * 5. Call prefetch + * 6. Ensure the correct number of materialize calls were made (one per aggregated range) + * 7. Verify the stored slices contain the correct ranges + */ + + // 1-2. Create a DummyLazyRead with 64 bytes + DummyLazyRead dummyLazyRead = createDummyLazyRead(64); + + // 3. Make a testable AggregatingSliceTrackingLazyRead + TestableAggregatingSliceTracker sliceTracking = new TestableAggregatingSliceTracker(dummyLazyRead); + + /* + * Non-adjacent ranges + */ + // 4. Create a list of ranges (two non-overlapping ranges with a gap) + List ranges = Arrays.asList( + Range.at(10, 5), // offset 10, length 5 (bytes 10-14) + Range.at(50, 10) // offset 50, length 10 (bytes 50-59) + ); + + // 5. Call prefetch + sliceTracking.prefetch(ranges); + + // 6. Ensure exactly 2 materialize calls were made + // AggregatingSliceTrackingLazyRead aggregates overlapping/adjacent ranges + // Since these ranges are not adjacent or overlapping, it makes 2 separate calls + assertEquals("AggregatingSliceTrackingLazyRead should make 2 materialize calls for non-adjacent ranges", + 2, dummyLazyRead.getNumMaterializeCalls()); + + // 7. Verify the stored slices contain two separate ranges + assertStoredSlices(sliceTracking, Arrays.asList( + Range.at(10, 5), // First slice + Range.at(50, 10) // Second slice + )); + + /* + * Adjacent ranges + */ + + // new sliceTracking instance to clear slices + sliceTracking = new TestableAggregatingSliceTracker(dummyLazyRead); + dummyLazyRead.resetNumMaterializeCalls(); + + // 4. Create a list of three contiguous ranges + List adjacentRanges = Arrays.asList( + Range.at(10, 5), // offset 10, length 5 (bytes 10-14) + Range.at(15, 10), // offset 15, length 10 (bytes 15-24) + Range.at(25, 5) // offset 25, length 5 (bytes 25-29) + ); + + // 5. Call prefetch + sliceTracking.prefetch(adjacentRanges); + + // 6. Ensure exactly 1 materialize call was made + // AggregatingSliceTrackingLazyRead should aggregate these three contiguous ranges + // into a single range from offset 10 to 30, with length 20 + assertEquals("AggregatingSliceTrackingLazyRead should make 1 materialize call for contiguous ranges", + 1, dummyLazyRead.getNumMaterializeCalls()); + + // 7. Verify the stored slices now contain three ranges total: + // the two from the first prefetch plus one aggregated range from the second prefetch + assertStoredSlices(sliceTracking, Arrays.asList( + Range.at(10, 20) // Aggregated range + )); + + } + + private static DummyLazyRead createDummyLazyRead(int size) { + byte[] data = new byte[size]; + for (int i = 0; i < data.length; i++) { + data[i] = (byte)i; + } + return new DummyLazyRead(ReadData.from(data)); + } + + /** + * Helper method to verify that stored slices match expected ranges. + * + * @param sliceTracking the SliceTrackingLazyRead instance + * @param expectedRanges the expected ranges stored in slices + */ + private static void assertStoredSlices(TestableSliceTracker sliceTracking, List expectedRanges) { + // Access protected slices field via a test helper + List actualSlices = sliceTracking.getSlices(); + + assertEquals("Number of stored slices should match", expectedRanges.size(), actualSlices.size()); + + for (int i = 0; i < expectedRanges.size(); i++) { + Range expected = expectedRanges.get(i); + Range actual = actualSlices.get(i); + assertEquals("Slice " + i + " offset should match", expected.offset(), actual.offset()); + assertEquals("Slice " + i + " length should match", expected.length(), actual.length()); + } + } + + /** + * Testable wrapper for DefaultSliceTrackingLazyRead that exposes slices. + */ + static class TestableDefaultSliceTracker extends DefaultSliceTrackingLazyRead implements TestableSliceTracker { + public TestableDefaultSliceTracker(LazyRead delegate) { + super(delegate); + } + + @Override + public List getSlices() { + return java.util.Collections.unmodifiableList(slices); + } + } + + /** + * Testable wrapper for AggregatingSliceTrackingLazyRead that exposes slices. + */ + static class TestableAggregatingSliceTracker extends AggregatingSliceTrackingLazyRead implements TestableSliceTracker { + public TestableAggregatingSliceTracker(LazyRead delegate) { + super(delegate); + } + + @Override + public List getSlices() { + return java.util.Collections.unmodifiableList(slices); + } + } + + /** + * Interface for testable slice trackers. + */ + interface TestableSliceTracker { + List getSlices(); + } + + static class DummyLazyRead implements LazyRead { + + private ReadData data; + private int numMaterializeCalls = 0; + + public DummyLazyRead( ReadData data ) { + this.data = data; + } + + @Override + public void close() throws IOException { + // no op + } + + @Override + public ReadData materialize(long offset, long length) throws N5IOException { + + numMaterializeCalls++; + return data.slice(offset, length).materialize(); + } + + @Override + public long size() throws N5IOException { + + return data.length(); + } + + public int getNumMaterializeCalls() { + + return numMaterializeCalls; + } + + public void resetNumMaterializeCalls() { + + numMaterializeCalls = 0; + } + + } +} From dc7907db9f1556ff7e7e365e4485e1f76b517dbe Mon Sep 17 00:00:00 2001 From: John Bogovic Date: Mon, 9 Feb 2026 17:27:41 -0500 Subject: [PATCH 05/10] wip: toward using prefetching --- .../n5/shard/DefaultDatasetAccess.java | 11 ++++++++++- .../janelia/saalfeldlab/n5/shard/RawShard.java | 17 +++++++++++++++++ 2 files changed, 27 insertions(+), 1 deletion(-) diff --git a/src/main/java/org/janelia/saalfeldlab/n5/shard/DefaultDatasetAccess.java b/src/main/java/org/janelia/saalfeldlab/n5/shard/DefaultDatasetAccess.java index 630c30033..66e0dd548 100644 --- a/src/main/java/org/janelia/saalfeldlab/n5/shard/DefaultDatasetAccess.java +++ b/src/main/java/org/janelia/saalfeldlab/n5/shard/DefaultDatasetAccess.java @@ -155,6 +155,15 @@ private void readChunksRecursive( // TODO: collect all the elementPos that we will need and prefetch // Probably best to add a prefetch method to RawShard? + // Here's an attempt at that. + // Don't love that we have to build a list of positions + // Consider making DataBlockRequest package private and passing them directly + final ArrayList positions = new ArrayList<>(); + for (final ChunkRequest request : requests) { + positions.add(request.position.relative(0)); + } + shard.prefetch(positions); + for (final ChunkRequest request : requests) { final long[] elementPos = request.position.relative(0); final ReadData elementData = shard.getElementData(elementPos); @@ -1015,7 +1024,7 @@ public List> chunks(final List> duplicates) { * Construct {@code ChunkRequests} from a list of level-0 grid positions * for reading. *

- * The nesting level ot the returned {@code ChunkRequests} is {@code + * The nesting level of the returned {@code ChunkRequests} is {@code * grid.numLevels()}, that is level of the highest-order shard + 1. This * implies that the requests are not guaranteed to be in the same shard (at * any level. {@link ChunkRequests#split() Splitting} the {@code diff --git a/src/main/java/org/janelia/saalfeldlab/n5/shard/RawShard.java b/src/main/java/org/janelia/saalfeldlab/n5/shard/RawShard.java index 01d7e961b..419ee1ccb 100644 --- a/src/main/java/org/janelia/saalfeldlab/n5/shard/RawShard.java +++ b/src/main/java/org/janelia/saalfeldlab/n5/shard/RawShard.java @@ -28,6 +28,11 @@ */ package org.janelia.saalfeldlab.n5.shard; +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; + +import org.janelia.saalfeldlab.n5.readdata.Range; import org.janelia.saalfeldlab.n5.readdata.ReadData; import org.janelia.saalfeldlab.n5.readdata.segment.Segment; import org.janelia.saalfeldlab.n5.readdata.segment.SegmentedReadData; @@ -86,4 +91,16 @@ public void setElementData(final ReadData data, final long[] pos) { final Segment segment = data == null ? null : SegmentedReadData.wrap(data).segments().get(0); index.set(segment, pos); } + + public void prefetch(List positions) { + + final List ranges = new ArrayList<>(positions.size()); + for (long[] pos : positions) { + final Segment seg = index.get(pos); + if (seg != null) + ranges.add(sourceData.location(seg)); + } + sourceData.prefetch(ranges); + } + } From 405e05061b2591340836e59d6309765f97e40b36 Mon Sep 17 00:00:00 2001 From: Caleb Hulbert Date: Thu, 26 Mar 2026 11:08:11 -0400 Subject: [PATCH 06/10] refactor: migrate previous batching/aggregate work to current development feat: use the aggregate/prefetch logic with VolatileReadData. By default, `VolatileReadData.from` now uses AggregatingSliceTrackingLazyRead. This is useful to ensure the default VolatileReadData behaves reasonably even for reading many chunks from the same shard. --- .../n5/readdata/VolatileReadData.java | 5 +- .../n5/backward/CompatibilityTest.java | 1 - .../n5/benchmarks/ReadDataBenchmarks.java | 3 - .../n5/http/HttpKeyValueAccessTest.java | 3 - .../n5/kva/TrackingKeyValueAccess.java | 24 ++++-- .../readdata/DelegatingVolatileReadData.java | 82 +++++++++++++++++++ .../n5/readdata/ReadDataTests.java | 1 - .../saalfeldlab/n5/shard/ShardTest.java | 36 +++++--- 8 files changed, 128 insertions(+), 27 deletions(-) create mode 100644 src/test/java/org/janelia/saalfeldlab/n5/readdata/DelegatingVolatileReadData.java diff --git a/src/main/java/org/janelia/saalfeldlab/n5/readdata/VolatileReadData.java b/src/main/java/org/janelia/saalfeldlab/n5/readdata/VolatileReadData.java index 122a399ef..7dba9e324 100644 --- a/src/main/java/org/janelia/saalfeldlab/n5/readdata/VolatileReadData.java +++ b/src/main/java/org/janelia/saalfeldlab/n5/readdata/VolatileReadData.java @@ -1,7 +1,7 @@ package org.janelia.saalfeldlab.n5.readdata; -import java.io.InputStream; import org.janelia.saalfeldlab.n5.N5Exception.N5IOException; +import org.janelia.saalfeldlab.n5.readdata.prefetch.AggregatingSliceTrackingLazyRead; /** * During its life-time, the content of a {@code VolatileReadData} should not be @@ -29,7 +29,8 @@ public interface VolatileReadData extends ReadData, AutoCloseable { * @return a new VolatileReadData */ static VolatileReadData from(final LazyRead lazyRead) { - return new LazyReadData(lazyRead); + final LazyRead aggregatingLazyRead = new AggregatingSliceTrackingLazyRead(lazyRead); + return new LazyReadData(aggregatingLazyRead); } } diff --git a/src/test/java/org/janelia/saalfeldlab/n5/backward/CompatibilityTest.java b/src/test/java/org/janelia/saalfeldlab/n5/backward/CompatibilityTest.java index e5a705422..17d421750 100644 --- a/src/test/java/org/janelia/saalfeldlab/n5/backward/CompatibilityTest.java +++ b/src/test/java/org/janelia/saalfeldlab/n5/backward/CompatibilityTest.java @@ -35,7 +35,6 @@ import java.io.File; import java.io.IOException; -import java.io.InputStream; import java.net.URI; import java.nio.file.Files; import java.util.Arrays; diff --git a/src/test/java/org/janelia/saalfeldlab/n5/benchmarks/ReadDataBenchmarks.java b/src/test/java/org/janelia/saalfeldlab/n5/benchmarks/ReadDataBenchmarks.java index 4ce148c50..73b6a2f46 100644 --- a/src/test/java/org/janelia/saalfeldlab/n5/benchmarks/ReadDataBenchmarks.java +++ b/src/test/java/org/janelia/saalfeldlab/n5/benchmarks/ReadDataBenchmarks.java @@ -29,8 +29,6 @@ package org.janelia.saalfeldlab.n5.benchmarks; import java.io.IOException; -import java.io.OutputStream; -import java.nio.file.FileSystems; import java.nio.file.Files; import java.nio.file.Path; import java.util.ArrayList; @@ -40,7 +38,6 @@ import org.janelia.saalfeldlab.n5.FileSystemKeyValueAccess; import org.janelia.saalfeldlab.n5.KeyValueAccess; -import org.janelia.saalfeldlab.n5.LockedChannel; import org.janelia.saalfeldlab.n5.N5Exception; import org.janelia.saalfeldlab.n5.readdata.ReadData; import org.janelia.saalfeldlab.n5.readdata.VolatileReadData; diff --git a/src/test/java/org/janelia/saalfeldlab/n5/http/HttpKeyValueAccessTest.java b/src/test/java/org/janelia/saalfeldlab/n5/http/HttpKeyValueAccessTest.java index 8572e1cee..fd151a265 100644 --- a/src/test/java/org/janelia/saalfeldlab/n5/http/HttpKeyValueAccessTest.java +++ b/src/test/java/org/janelia/saalfeldlab/n5/http/HttpKeyValueAccessTest.java @@ -30,17 +30,14 @@ import org.apache.commons.io.IOUtils; import org.janelia.saalfeldlab.n5.HttpKeyValueAccess; -import org.janelia.saalfeldlab.n5.LockedChannel; import org.janelia.saalfeldlab.n5.N5Exception; import org.janelia.saalfeldlab.n5.readdata.ReadData; import org.janelia.saalfeldlab.n5.readdata.VolatileReadData; import org.junit.Test; import java.io.IOException; -import java.io.InputStream; import java.net.URI; import java.nio.charset.Charset; -import java.util.function.Function; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertThrows; diff --git a/src/test/java/org/janelia/saalfeldlab/n5/kva/TrackingKeyValueAccess.java b/src/test/java/org/janelia/saalfeldlab/n5/kva/TrackingKeyValueAccess.java index 0c4f4a876..a72608363 100644 --- a/src/test/java/org/janelia/saalfeldlab/n5/kva/TrackingKeyValueAccess.java +++ b/src/test/java/org/janelia/saalfeldlab/n5/kva/TrackingKeyValueAccess.java @@ -2,16 +2,18 @@ import org.janelia.saalfeldlab.n5.KeyValueAccess; import org.janelia.saalfeldlab.n5.N5Exception; +import org.janelia.saalfeldlab.n5.readdata.DelegatingVolatileReadData; import org.janelia.saalfeldlab.n5.readdata.LazyRead; import org.janelia.saalfeldlab.n5.readdata.ReadData; import org.janelia.saalfeldlab.n5.readdata.VolatileReadData; -import org.janelia.saalfeldlab.n5.shard.ShardTest; +import org.janelia.saalfeldlab.n5.readdata.prefetch.AggregatingSliceTrackingLazyRead; public class TrackingKeyValueAccess extends DelegateKeyValueAccess { public int numMaterializeCalls = 0; public int numIsFileCalls = 0; public long totalBytesRead = 0; + public boolean aggregate = false; public TrackingKeyValueAccess(final KeyValueAccess kva) { super(kva); @@ -25,15 +27,27 @@ public boolean isFile(String normalPath) { @Override public VolatileReadData createReadData(final String normalPath) { -// throw new N5NoSuchKeyException("Test No Such Key"); - return VolatileReadData.from(new TrackingVolatileReadData(kva.createReadData(normalPath))); + + final VolatileReadData volatileReadData = kva.createReadData(normalPath); + final TrackingLazyRead trackingLazyRead = new TrackingLazyRead(volatileReadData); + LazyRead lazyRead = trackingLazyRead; + if (aggregate) + lazyRead = new AggregatingSliceTrackingLazyRead(trackingLazyRead); + VolatileReadData delegate = VolatileReadData.from( lazyRead ); + return new DelegatingVolatileReadData(delegate) { + @Override + public void close() throws N5Exception.N5IOException { + super.close(); //closes delegate + volatileReadData.close(); + } + }; } - private class TrackingVolatileReadData implements LazyRead { + private class TrackingLazyRead implements LazyRead { private final VolatileReadData readData; - TrackingVolatileReadData(final VolatileReadData readData) { + TrackingLazyRead(final VolatileReadData readData) { this.readData = readData; } diff --git a/src/test/java/org/janelia/saalfeldlab/n5/readdata/DelegatingVolatileReadData.java b/src/test/java/org/janelia/saalfeldlab/n5/readdata/DelegatingVolatileReadData.java new file mode 100644 index 000000000..d2f633f1c --- /dev/null +++ b/src/test/java/org/janelia/saalfeldlab/n5/readdata/DelegatingVolatileReadData.java @@ -0,0 +1,82 @@ +package org.janelia.saalfeldlab.n5.readdata; + +import org.janelia.saalfeldlab.n5.N5Exception; + +import java.io.InputStream; +import java.io.OutputStream; +import java.nio.ByteBuffer; +import java.util.Collection; + +public class DelegatingVolatileReadData implements VolatileReadData { + + private final VolatileReadData delegate; + + public DelegatingVolatileReadData(VolatileReadData delegate) { + this.delegate = delegate; + } + + @Override + public void close() throws N5Exception.N5IOException { + delegate.close(); + } + + @Override + public long length() { + return delegate.length(); + } + + @Override + public long requireLength() throws N5Exception.N5IOException { + return delegate.requireLength(); + } + + @Override + public ReadData limit(long length) throws N5Exception.N5IOException { + return delegate.limit(length); + } + + @Override + public ReadData slice(long offset, long length) throws N5Exception.N5IOException { + return delegate.slice(offset, length); + } + + @Override + public ReadData slice(Range range) throws N5Exception.N5IOException { + return delegate.slice(range); + } + + @Override + public InputStream inputStream() throws N5Exception.N5IOException, IllegalStateException { + return delegate.inputStream(); + } + + @Override + public byte[] allBytes() throws N5Exception.N5IOException, IllegalStateException { + return delegate.allBytes(); + } + + @Override + public ByteBuffer toByteBuffer() throws N5Exception.N5IOException, IllegalStateException { + return delegate.toByteBuffer(); + } + + @Override + public ReadData materialize() throws N5Exception.N5IOException { + return delegate.materialize(); + } + + @Override + public void writeTo(OutputStream outputStream) throws N5Exception.N5IOException, IllegalStateException { + delegate.writeTo(outputStream); + } + + @Override + public void prefetch(Collection ranges) throws N5Exception.N5IOException { + delegate.prefetch(ranges); + } + + @Override + public ReadData encode(OutputStreamOperator encoder) { + return delegate.encode(encoder); + } +} diff --git a/src/test/java/org/janelia/saalfeldlab/n5/readdata/ReadDataTests.java b/src/test/java/org/janelia/saalfeldlab/n5/readdata/ReadDataTests.java index c4f3fa51b..78522ff3e 100644 --- a/src/test/java/org/janelia/saalfeldlab/n5/readdata/ReadDataTests.java +++ b/src/test/java/org/janelia/saalfeldlab/n5/readdata/ReadDataTests.java @@ -38,7 +38,6 @@ import java.io.IOException; import java.io.InputStream; import java.io.OutputStream; -import java.nio.file.FileSystems; import java.util.Arrays; import java.util.function.IntUnaryOperator; diff --git a/src/test/java/org/janelia/saalfeldlab/n5/shard/ShardTest.java b/src/test/java/org/janelia/saalfeldlab/n5/shard/ShardTest.java index 7d30db739..36bf19240 100644 --- a/src/test/java/org/janelia/saalfeldlab/n5/shard/ShardTest.java +++ b/src/test/java/org/janelia/saalfeldlab/n5/shard/ShardTest.java @@ -537,9 +537,8 @@ public void numReadsTest() { new ByteArrayDataBlock(chunkSize, new long[]{11, 11}, data) ); - writer.resetNumMaterializeCalls(); - writer.readChunks(dataset, datasetAttributes, Collections.singletonList(new long[] {0,0})); - System.out.println(writer.getNumMaterializeCalls()); + writer.resetNumMaterializeCalls(); + writer.readChunks(dataset, datasetAttributes, Collections.singletonList(new long[] {0,0})); ArrayList ptList = new ArrayList<>(); ptList.add(new long[] {0, 0}); @@ -547,10 +546,8 @@ public void numReadsTest() { ptList.add(new long[] {1, 0}); ptList.add(new long[] {1, 1}); - writer.resetNumMaterializeCalls(); - writer.readChunks(dataset, datasetAttributes, ptList); - System.out.println(writer.getNumMaterializeCalls()); - System.out.println(""); + writer.resetNumMaterializeCalls(); + writer.readChunks(dataset, datasetAttributes, ptList); } @Test @@ -656,14 +653,29 @@ public void testPartialReadAggregationBehavior() { writer.resetNumMaterializeCalls(); writer.readChunks(dataset, datasetAttributes, ptList); - // TODO change this if and when we implement aggregation of read calls - // one for the index, one for each of the four blocks - assertEquals(5, writer.getNumMaterializeCalls()); + // one for the index, one for the four blocks (aggregated) + assertEquals(2, writer.getNumMaterializeCalls()); writer.resetNumMaterializeCalls(); writer.readBlock(dataset, datasetAttributes, new long[] {0,0}); - // one for the index, one for each of the four blocks - assertEquals(5, writer.getNumMaterializeCalls()); + // one for the index, one for the four blocks (aggregated) + assertEquals(2, writer.getNumMaterializeCalls()); + + + /** + * Aggregate read calls + */ + writer.tkva.aggregate = true; + writer.resetNumMaterializeCalls(); + writer.readChunks(dataset, datasetAttributes, ptList); + + // one for the index, one that covers ALL the blocks) + assertEquals(2, writer.getNumMaterializeCalls()); + + writer.resetNumMaterializeCalls(); + writer.readBlock(dataset, datasetAttributes, new long[] {0,0}); + // one for the index, one that covers ALL the blocks + assertEquals(2, writer.getNumMaterializeCalls()); } } From 03e96693d777283b8a5d55841b5a519ea9993ecb Mon Sep 17 00:00:00 2001 From: tpietzsch Date: Fri, 3 Apr 2026 22:45:07 +0200 Subject: [PATCH 07/10] clean up The wrapped volatileReadData is closed via the wrapper chain LazyReadData -> AggregatingSliceTrackingLazyRead -> TrackingLazyRead -> volatileReadData No need for DelegatingVolatileReadData. --- .../n5/kva/TrackingKeyValueAccess.java | 10 +-- .../readdata/DelegatingVolatileReadData.java | 82 ------------------- 2 files changed, 1 insertion(+), 91 deletions(-) delete mode 100644 src/test/java/org/janelia/saalfeldlab/n5/readdata/DelegatingVolatileReadData.java diff --git a/src/test/java/org/janelia/saalfeldlab/n5/kva/TrackingKeyValueAccess.java b/src/test/java/org/janelia/saalfeldlab/n5/kva/TrackingKeyValueAccess.java index a72608363..3f747b8ec 100644 --- a/src/test/java/org/janelia/saalfeldlab/n5/kva/TrackingKeyValueAccess.java +++ b/src/test/java/org/janelia/saalfeldlab/n5/kva/TrackingKeyValueAccess.java @@ -2,7 +2,6 @@ import org.janelia.saalfeldlab.n5.KeyValueAccess; import org.janelia.saalfeldlab.n5.N5Exception; -import org.janelia.saalfeldlab.n5.readdata.DelegatingVolatileReadData; import org.janelia.saalfeldlab.n5.readdata.LazyRead; import org.janelia.saalfeldlab.n5.readdata.ReadData; import org.janelia.saalfeldlab.n5.readdata.VolatileReadData; @@ -33,14 +32,7 @@ public VolatileReadData createReadData(final String normalPath) { LazyRead lazyRead = trackingLazyRead; if (aggregate) lazyRead = new AggregatingSliceTrackingLazyRead(trackingLazyRead); - VolatileReadData delegate = VolatileReadData.from( lazyRead ); - return new DelegatingVolatileReadData(delegate) { - @Override - public void close() throws N5Exception.N5IOException { - super.close(); //closes delegate - volatileReadData.close(); - } - }; + return VolatileReadData.from( lazyRead ); } private class TrackingLazyRead implements LazyRead { diff --git a/src/test/java/org/janelia/saalfeldlab/n5/readdata/DelegatingVolatileReadData.java b/src/test/java/org/janelia/saalfeldlab/n5/readdata/DelegatingVolatileReadData.java deleted file mode 100644 index d2f633f1c..000000000 --- a/src/test/java/org/janelia/saalfeldlab/n5/readdata/DelegatingVolatileReadData.java +++ /dev/null @@ -1,82 +0,0 @@ -package org.janelia.saalfeldlab.n5.readdata; - -import org.janelia.saalfeldlab.n5.N5Exception; - -import java.io.InputStream; -import java.io.OutputStream; -import java.nio.ByteBuffer; -import java.util.Collection; - -public class DelegatingVolatileReadData implements VolatileReadData { - - private final VolatileReadData delegate; - - public DelegatingVolatileReadData(VolatileReadData delegate) { - this.delegate = delegate; - } - - @Override - public void close() throws N5Exception.N5IOException { - delegate.close(); - } - - @Override - public long length() { - return delegate.length(); - } - - @Override - public long requireLength() throws N5Exception.N5IOException { - return delegate.requireLength(); - } - - @Override - public ReadData limit(long length) throws N5Exception.N5IOException { - return delegate.limit(length); - } - - @Override - public ReadData slice(long offset, long length) throws N5Exception.N5IOException { - return delegate.slice(offset, length); - } - - @Override - public ReadData slice(Range range) throws N5Exception.N5IOException { - return delegate.slice(range); - } - - @Override - public InputStream inputStream() throws N5Exception.N5IOException, IllegalStateException { - return delegate.inputStream(); - } - - @Override - public byte[] allBytes() throws N5Exception.N5IOException, IllegalStateException { - return delegate.allBytes(); - } - - @Override - public ByteBuffer toByteBuffer() throws N5Exception.N5IOException, IllegalStateException { - return delegate.toByteBuffer(); - } - - @Override - public ReadData materialize() throws N5Exception.N5IOException { - return delegate.materialize(); - } - - @Override - public void writeTo(OutputStream outputStream) throws N5Exception.N5IOException, IllegalStateException { - delegate.writeTo(outputStream); - } - - @Override - public void prefetch(Collection ranges) throws N5Exception.N5IOException { - delegate.prefetch(ranges); - } - - @Override - public ReadData encode(OutputStreamOperator encoder) { - return delegate.encode(encoder); - } -} From 8dec89a862c34591f6b8410311e5ce7b2998b79f Mon Sep 17 00:00:00 2001 From: tpietzsch Date: Fri, 3 Apr 2026 23:02:15 +0200 Subject: [PATCH 08/10] Make SliceTrackingLazyRead non-abstract It doesn't do any prefetching, just keeps track of what has been materialized (and uses that to avoid materializing slices that are already fully covered). --- .../DefaultSliceTrackingLazyRead.java | 5 -- .../prefetch/SliceTrackingLazyRead.java | 46 +++++-------------- 2 files changed, 11 insertions(+), 40 deletions(-) diff --git a/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/DefaultSliceTrackingLazyRead.java b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/DefaultSliceTrackingLazyRead.java index 1e607ddc1..f35a8733c 100644 --- a/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/DefaultSliceTrackingLazyRead.java +++ b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/DefaultSliceTrackingLazyRead.java @@ -42,9 +42,4 @@ public void prefetch(final Collection ranges) throws N5IOExcept materialize(fromIndex, toIndex - fromIndex); } } - - private boolean isCovered(final Range slice) { - - return Slices.findContainingSlice(slices, slice) != null; - } } diff --git a/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyRead.java b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyRead.java index de20f2c39..5835c473e 100644 --- a/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyRead.java +++ b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyRead.java @@ -2,14 +2,22 @@ import java.io.IOException; import java.util.ArrayList; -import java.util.Collection; import java.util.List; import org.janelia.saalfeldlab.n5.N5Exception.N5IOException; import org.janelia.saalfeldlab.n5.readdata.LazyRead; import org.janelia.saalfeldlab.n5.readdata.Range; import org.janelia.saalfeldlab.n5.readdata.ReadData; -public abstract class SliceTrackingLazyRead implements LazyRead { +/** + * A {@link LazyRead} that wraps a delegate {@code LazyRead} and keeps track of + * all slices that have been {@link #materialize materialized}. + *

+ * When materializing a new slice, we first check whether it is completely + * covered by a materialized slice that we already track. If so, then we just + * return a slice on the existing materialized slice. If not, we materialize the + * slice from the delegate track it. + */ +public class SliceTrackingLazyRead implements LazyRead { protected static class Slice implements Range { @@ -75,39 +83,7 @@ public long size() throws N5IOException { return delegate.size(); } - /** - * Indicates that the given slices will be subsequently read. - * {@code LazyRead} implementations (optionally) may take steps to prepare - * for these subsequent slices. - *

- * Minimal implementation: Find offset and length covering all ranges that - * are not yet fully covered by existing slices. Then materialize the slice - * covering that range. - * - * @param ranges - * slice ranges to prefetch - * - * @throws N5IOException - * if any I/O error occurs - */ - @Override - public void prefetch(final Collection ranges) throws N5IOException { - - long fromIndex = Long.MAX_VALUE; - long toIndex = Long.MIN_VALUE; - for (final Range slice : ranges) { - if (!isCovered(slice)) { - fromIndex = Math.min(fromIndex, slice.offset()); - toIndex = Math.max(toIndex, slice.end()); - } - } - - if (fromIndex < toIndex) { - materialize(fromIndex, toIndex - fromIndex); - } - } - - private boolean isCovered(final Range slice) { + protected boolean isCovered(final Range slice) { return Slices.findContainingSlice(slices, slice) != null; } From 0a0e33a008bf6b309ebbabec793ae55d3476b18f Mon Sep 17 00:00:00 2001 From: tpietzsch Date: Fri, 3 Apr 2026 23:12:47 +0200 Subject: [PATCH 09/10] rename prefetching LazyReads to "...PrefetchLazyRead" instead of "...SliceTracking..." which is not their main feature --- .../saalfeldlab/n5/readdata/VolatileReadData.java | 4 ++-- ...ingLazyRead.java => AggregatingPrefetchLazyRead.java} | 9 +++++++-- ...ckingLazyRead.java => EnclosingPrefetchLazyRead.java} | 8 ++++++-- .../saalfeldlab/n5/kva/TrackingKeyValueAccess.java | 4 ++-- .../n5/readdata/prefetch/SliceTrackingLazyReadTests.java | 4 ++-- 5 files changed, 19 insertions(+), 10 deletions(-) rename src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/{AggregatingSliceTrackingLazyRead.java => AggregatingPrefetchLazyRead.java} (73%) rename src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/{DefaultSliceTrackingLazyRead.java => EnclosingPrefetchLazyRead.java} (81%) diff --git a/src/main/java/org/janelia/saalfeldlab/n5/readdata/VolatileReadData.java b/src/main/java/org/janelia/saalfeldlab/n5/readdata/VolatileReadData.java index 7dba9e324..25da8b0b8 100644 --- a/src/main/java/org/janelia/saalfeldlab/n5/readdata/VolatileReadData.java +++ b/src/main/java/org/janelia/saalfeldlab/n5/readdata/VolatileReadData.java @@ -1,7 +1,7 @@ package org.janelia.saalfeldlab.n5.readdata; import org.janelia.saalfeldlab.n5.N5Exception.N5IOException; -import org.janelia.saalfeldlab.n5.readdata.prefetch.AggregatingSliceTrackingLazyRead; +import org.janelia.saalfeldlab.n5.readdata.prefetch.AggregatingPrefetchLazyRead; /** * During its life-time, the content of a {@code VolatileReadData} should not be @@ -29,7 +29,7 @@ public interface VolatileReadData extends ReadData, AutoCloseable { * @return a new VolatileReadData */ static VolatileReadData from(final LazyRead lazyRead) { - final LazyRead aggregatingLazyRead = new AggregatingSliceTrackingLazyRead(lazyRead); + final LazyRead aggregatingLazyRead = new AggregatingPrefetchLazyRead(lazyRead); return new LazyReadData(aggregatingLazyRead); } diff --git a/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/AggregatingSliceTrackingLazyRead.java b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/AggregatingPrefetchLazyRead.java similarity index 73% rename from src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/AggregatingSliceTrackingLazyRead.java rename to src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/AggregatingPrefetchLazyRead.java index fc4fb660d..3a5145359 100644 --- a/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/AggregatingSliceTrackingLazyRead.java +++ b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/AggregatingPrefetchLazyRead.java @@ -6,9 +6,14 @@ import org.janelia.saalfeldlab.n5.readdata.LazyRead; import org.janelia.saalfeldlab.n5.readdata.Range; -public class AggregatingSliceTrackingLazyRead extends SliceTrackingLazyRead { +/** + * A {@link SliceTrackingLazyRead} that implements {@link #prefetch} to + * aggregate overlapping / adjacent ranges and then materialize each aggregated + * range. + */ +public class AggregatingPrefetchLazyRead extends SliceTrackingLazyRead { - public AggregatingSliceTrackingLazyRead(final LazyRead delegate) { + public AggregatingPrefetchLazyRead(final LazyRead delegate) { super(delegate); } diff --git a/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/DefaultSliceTrackingLazyRead.java b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/EnclosingPrefetchLazyRead.java similarity index 81% rename from src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/DefaultSliceTrackingLazyRead.java rename to src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/EnclosingPrefetchLazyRead.java index f35a8733c..b0d6fdc7a 100644 --- a/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/DefaultSliceTrackingLazyRead.java +++ b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/EnclosingPrefetchLazyRead.java @@ -5,9 +5,13 @@ import org.janelia.saalfeldlab.n5.readdata.LazyRead; import org.janelia.saalfeldlab.n5.readdata.Range; -public class DefaultSliceTrackingLazyRead extends SliceTrackingLazyRead { +/** + * A {@link SliceTrackingLazyRead} that implements {@link #prefetch} to + * materialize the bounding range of all requested ranges. + */ +public class EnclosingPrefetchLazyRead extends SliceTrackingLazyRead { - public DefaultSliceTrackingLazyRead(final LazyRead delegate) { + public EnclosingPrefetchLazyRead(final LazyRead delegate) { super(delegate); } diff --git a/src/test/java/org/janelia/saalfeldlab/n5/kva/TrackingKeyValueAccess.java b/src/test/java/org/janelia/saalfeldlab/n5/kva/TrackingKeyValueAccess.java index 3f747b8ec..a2e18b17d 100644 --- a/src/test/java/org/janelia/saalfeldlab/n5/kva/TrackingKeyValueAccess.java +++ b/src/test/java/org/janelia/saalfeldlab/n5/kva/TrackingKeyValueAccess.java @@ -5,7 +5,7 @@ import org.janelia.saalfeldlab.n5.readdata.LazyRead; import org.janelia.saalfeldlab.n5.readdata.ReadData; import org.janelia.saalfeldlab.n5.readdata.VolatileReadData; -import org.janelia.saalfeldlab.n5.readdata.prefetch.AggregatingSliceTrackingLazyRead; +import org.janelia.saalfeldlab.n5.readdata.prefetch.AggregatingPrefetchLazyRead; public class TrackingKeyValueAccess extends DelegateKeyValueAccess { @@ -31,7 +31,7 @@ public VolatileReadData createReadData(final String normalPath) { final TrackingLazyRead trackingLazyRead = new TrackingLazyRead(volatileReadData); LazyRead lazyRead = trackingLazyRead; if (aggregate) - lazyRead = new AggregatingSliceTrackingLazyRead(trackingLazyRead); + lazyRead = new AggregatingPrefetchLazyRead(trackingLazyRead); return VolatileReadData.from( lazyRead ); } diff --git a/src/test/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyReadTests.java b/src/test/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyReadTests.java index 044e6c1a3..fcac9a5e3 100644 --- a/src/test/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyReadTests.java +++ b/src/test/java/org/janelia/saalfeldlab/n5/readdata/prefetch/SliceTrackingLazyReadTests.java @@ -160,7 +160,7 @@ private static void assertStoredSlices(TestableSliceTracker sliceTracking, List< /** * Testable wrapper for DefaultSliceTrackingLazyRead that exposes slices. */ - static class TestableDefaultSliceTracker extends DefaultSliceTrackingLazyRead implements TestableSliceTracker { + static class TestableDefaultSliceTracker extends EnclosingPrefetchLazyRead implements TestableSliceTracker { public TestableDefaultSliceTracker(LazyRead delegate) { super(delegate); } @@ -174,7 +174,7 @@ public List getSlices() { /** * Testable wrapper for AggregatingSliceTrackingLazyRead that exposes slices. */ - static class TestableAggregatingSliceTracker extends AggregatingSliceTrackingLazyRead implements TestableSliceTracker { + static class TestableAggregatingSliceTracker extends AggregatingPrefetchLazyRead implements TestableSliceTracker { public TestableAggregatingSliceTracker(LazyRead delegate) { super(delegate); } From 6b55999a6be967eac8d398bdfe1e9545a6ec3d63 Mon Sep 17 00:00:00 2001 From: tpietzsch Date: Fri, 3 Apr 2026 23:18:37 +0200 Subject: [PATCH 10/10] Don't include ranges that are already materialized in the aggregation --- .../n5/readdata/prefetch/AggregatingPrefetchLazyRead.java | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/AggregatingPrefetchLazyRead.java b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/AggregatingPrefetchLazyRead.java index 3a5145359..dcc0097a1 100644 --- a/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/AggregatingPrefetchLazyRead.java +++ b/src/main/java/org/janelia/saalfeldlab/n5/readdata/prefetch/AggregatingPrefetchLazyRead.java @@ -1,7 +1,9 @@ package org.janelia.saalfeldlab.n5.readdata.prefetch; +import java.util.ArrayList; import java.util.Collection; +import java.util.List; import org.janelia.saalfeldlab.n5.N5Exception.N5IOException; import org.janelia.saalfeldlab.n5.readdata.LazyRead; import org.janelia.saalfeldlab.n5.readdata.Range; @@ -31,7 +33,9 @@ public AggregatingPrefetchLazyRead(final LazyRead delegate) { @Override public void prefetch(final Collection ranges) throws N5IOException { - final Collection aggregatedRanges = Range.aggregate(ranges); + final List filteredRanges = new ArrayList<>(ranges); + filteredRanges.removeIf(this::isCovered); + final Collection aggregatedRanges = Range.aggregate(filteredRanges); for (final Range slice : aggregatedRanges) { materialize(slice.offset(), slice.length()); }