Skip to content

Commit 405e050

Browse files
cmhulberttpietzsch
authored andcommitted
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.
1 parent dc7907d commit 405e050

8 files changed

Lines changed: 128 additions & 27 deletions

File tree

src/main/java/org/janelia/saalfeldlab/n5/readdata/VolatileReadData.java

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
package org.janelia.saalfeldlab.n5.readdata;
22

3-
import java.io.InputStream;
43
import org.janelia.saalfeldlab.n5.N5Exception.N5IOException;
4+
import org.janelia.saalfeldlab.n5.readdata.prefetch.AggregatingSliceTrackingLazyRead;
55

66
/**
77
* During its life-time, the content of a {@code VolatileReadData} should not be
@@ -29,7 +29,8 @@ public interface VolatileReadData extends ReadData, AutoCloseable {
2929
* @return a new VolatileReadData
3030
*/
3131
static VolatileReadData from(final LazyRead lazyRead) {
32-
return new LazyReadData(lazyRead);
32+
final LazyRead aggregatingLazyRead = new AggregatingSliceTrackingLazyRead(lazyRead);
33+
return new LazyReadData(aggregatingLazyRead);
3334
}
3435

3536
}

src/test/java/org/janelia/saalfeldlab/n5/backward/CompatibilityTest.java

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -35,7 +35,6 @@
3535

3636
import java.io.File;
3737
import java.io.IOException;
38-
import java.io.InputStream;
3938
import java.net.URI;
4039
import java.nio.file.Files;
4140
import java.util.Arrays;

src/test/java/org/janelia/saalfeldlab/n5/benchmarks/ReadDataBenchmarks.java

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -29,8 +29,6 @@
2929
package org.janelia.saalfeldlab.n5.benchmarks;
3030

3131
import java.io.IOException;
32-
import java.io.OutputStream;
33-
import java.nio.file.FileSystems;
3432
import java.nio.file.Files;
3533
import java.nio.file.Path;
3634
import java.util.ArrayList;
@@ -40,7 +38,6 @@
4038

4139
import org.janelia.saalfeldlab.n5.FileSystemKeyValueAccess;
4240
import org.janelia.saalfeldlab.n5.KeyValueAccess;
43-
import org.janelia.saalfeldlab.n5.LockedChannel;
4441
import org.janelia.saalfeldlab.n5.N5Exception;
4542
import org.janelia.saalfeldlab.n5.readdata.ReadData;
4643
import org.janelia.saalfeldlab.n5.readdata.VolatileReadData;

src/test/java/org/janelia/saalfeldlab/n5/http/HttpKeyValueAccessTest.java

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -30,17 +30,14 @@
3030

3131
import org.apache.commons.io.IOUtils;
3232
import org.janelia.saalfeldlab.n5.HttpKeyValueAccess;
33-
import org.janelia.saalfeldlab.n5.LockedChannel;
3433
import org.janelia.saalfeldlab.n5.N5Exception;
3534
import org.janelia.saalfeldlab.n5.readdata.ReadData;
3635
import org.janelia.saalfeldlab.n5.readdata.VolatileReadData;
3736
import org.junit.Test;
3837

3938
import java.io.IOException;
40-
import java.io.InputStream;
4139
import java.net.URI;
4240
import java.nio.charset.Charset;
43-
import java.util.function.Function;
4441

4542
import static org.junit.Assert.assertEquals;
4643
import static org.junit.Assert.assertThrows;

src/test/java/org/janelia/saalfeldlab/n5/kva/TrackingKeyValueAccess.java

Lines changed: 19 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -2,16 +2,18 @@
22

33
import org.janelia.saalfeldlab.n5.KeyValueAccess;
44
import org.janelia.saalfeldlab.n5.N5Exception;
5+
import org.janelia.saalfeldlab.n5.readdata.DelegatingVolatileReadData;
56
import org.janelia.saalfeldlab.n5.readdata.LazyRead;
67
import org.janelia.saalfeldlab.n5.readdata.ReadData;
78
import org.janelia.saalfeldlab.n5.readdata.VolatileReadData;
8-
import org.janelia.saalfeldlab.n5.shard.ShardTest;
9+
import org.janelia.saalfeldlab.n5.readdata.prefetch.AggregatingSliceTrackingLazyRead;
910

1011
public class TrackingKeyValueAccess extends DelegateKeyValueAccess {
1112

1213
public int numMaterializeCalls = 0;
1314
public int numIsFileCalls = 0;
1415
public long totalBytesRead = 0;
16+
public boolean aggregate = false;
1517

1618
public TrackingKeyValueAccess(final KeyValueAccess kva) {
1719
super(kva);
@@ -25,15 +27,27 @@ public boolean isFile(String normalPath) {
2527

2628
@Override
2729
public VolatileReadData createReadData(final String normalPath) {
28-
// throw new N5NoSuchKeyException("Test No Such Key");
29-
return VolatileReadData.from(new TrackingVolatileReadData(kva.createReadData(normalPath)));
30+
31+
final VolatileReadData volatileReadData = kva.createReadData(normalPath);
32+
final TrackingLazyRead trackingLazyRead = new TrackingLazyRead(volatileReadData);
33+
LazyRead lazyRead = trackingLazyRead;
34+
if (aggregate)
35+
lazyRead = new AggregatingSliceTrackingLazyRead(trackingLazyRead);
36+
VolatileReadData delegate = VolatileReadData.from( lazyRead );
37+
return new DelegatingVolatileReadData(delegate) {
38+
@Override
39+
public void close() throws N5Exception.N5IOException {
40+
super.close(); //closes delegate
41+
volatileReadData.close();
42+
}
43+
};
3044
}
3145

32-
private class TrackingVolatileReadData implements LazyRead {
46+
private class TrackingLazyRead implements LazyRead {
3347

3448
private final VolatileReadData readData;
3549

36-
TrackingVolatileReadData(final VolatileReadData readData) {
50+
TrackingLazyRead(final VolatileReadData readData) {
3751
this.readData = readData;
3852
}
3953

Lines changed: 82 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,82 @@
1+
package org.janelia.saalfeldlab.n5.readdata;
2+
3+
import org.janelia.saalfeldlab.n5.N5Exception;
4+
5+
import java.io.InputStream;
6+
import java.io.OutputStream;
7+
import java.nio.ByteBuffer;
8+
import java.util.Collection;
9+
10+
public class DelegatingVolatileReadData implements VolatileReadData {
11+
12+
private final VolatileReadData delegate;
13+
14+
public DelegatingVolatileReadData(VolatileReadData delegate) {
15+
this.delegate = delegate;
16+
}
17+
18+
@Override
19+
public void close() throws N5Exception.N5IOException {
20+
delegate.close();
21+
}
22+
23+
@Override
24+
public long length() {
25+
return delegate.length();
26+
}
27+
28+
@Override
29+
public long requireLength() throws N5Exception.N5IOException {
30+
return delegate.requireLength();
31+
}
32+
33+
@Override
34+
public ReadData limit(long length) throws N5Exception.N5IOException {
35+
return delegate.limit(length);
36+
}
37+
38+
@Override
39+
public ReadData slice(long offset, long length) throws N5Exception.N5IOException {
40+
return delegate.slice(offset, length);
41+
}
42+
43+
@Override
44+
public ReadData slice(Range range) throws N5Exception.N5IOException {
45+
return delegate.slice(range);
46+
}
47+
48+
@Override
49+
public InputStream inputStream() throws N5Exception.N5IOException, IllegalStateException {
50+
return delegate.inputStream();
51+
}
52+
53+
@Override
54+
public byte[] allBytes() throws N5Exception.N5IOException, IllegalStateException {
55+
return delegate.allBytes();
56+
}
57+
58+
@Override
59+
public ByteBuffer toByteBuffer() throws N5Exception.N5IOException, IllegalStateException {
60+
return delegate.toByteBuffer();
61+
}
62+
63+
@Override
64+
public ReadData materialize() throws N5Exception.N5IOException {
65+
return delegate.materialize();
66+
}
67+
68+
@Override
69+
public void writeTo(OutputStream outputStream) throws N5Exception.N5IOException, IllegalStateException {
70+
delegate.writeTo(outputStream);
71+
}
72+
73+
@Override
74+
public void prefetch(Collection<? extends Range> ranges) throws N5Exception.N5IOException {
75+
delegate.prefetch(ranges);
76+
}
77+
78+
@Override
79+
public ReadData encode(OutputStreamOperator encoder) {
80+
return delegate.encode(encoder);
81+
}
82+
}

src/test/java/org/janelia/saalfeldlab/n5/readdata/ReadDataTests.java

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -38,7 +38,6 @@
3838
import java.io.IOException;
3939
import java.io.InputStream;
4040
import java.io.OutputStream;
41-
import java.nio.file.FileSystems;
4241
import java.util.Arrays;
4342
import java.util.function.IntUnaryOperator;
4443

src/test/java/org/janelia/saalfeldlab/n5/shard/ShardTest.java

Lines changed: 24 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -537,20 +537,17 @@ public void numReadsTest() {
537537
new ByteArrayDataBlock(chunkSize, new long[]{11, 11}, data)
538538
);
539539

540-
writer.resetNumMaterializeCalls();
541-
writer.readChunks(dataset, datasetAttributes, Collections.singletonList(new long[] {0,0}));
542-
System.out.println(writer.getNumMaterializeCalls());
540+
writer.resetNumMaterializeCalls();
541+
writer.readChunks(dataset, datasetAttributes, Collections.singletonList(new long[] {0,0}));
543542

544543
ArrayList<long[]> ptList = new ArrayList<>();
545544
ptList.add(new long[] {0, 0});
546545
ptList.add(new long[] {0, 1});
547546
ptList.add(new long[] {1, 0});
548547
ptList.add(new long[] {1, 1});
549548

550-
writer.resetNumMaterializeCalls();
551-
writer.readChunks(dataset, datasetAttributes, ptList);
552-
System.out.println(writer.getNumMaterializeCalls());
553-
System.out.println("");
549+
writer.resetNumMaterializeCalls();
550+
writer.readChunks(dataset, datasetAttributes, ptList);
554551
}
555552

556553
@Test
@@ -656,14 +653,29 @@ public void testPartialReadAggregationBehavior() {
656653
writer.resetNumMaterializeCalls();
657654
writer.readChunks(dataset, datasetAttributes, ptList);
658655

659-
// TODO change this if and when we implement aggregation of read calls
660-
// one for the index, one for each of the four blocks
661-
assertEquals(5, writer.getNumMaterializeCalls());
656+
// one for the index, one for the four blocks (aggregated)
657+
assertEquals(2, writer.getNumMaterializeCalls());
662658

663659
writer.resetNumMaterializeCalls();
664660
writer.readBlock(dataset, datasetAttributes, new long[] {0,0});
665-
// one for the index, one for each of the four blocks
666-
assertEquals(5, writer.getNumMaterializeCalls());
661+
// one for the index, one for the four blocks (aggregated)
662+
assertEquals(2, writer.getNumMaterializeCalls());
663+
664+
665+
/**
666+
* Aggregate read calls
667+
*/
668+
writer.tkva.aggregate = true;
669+
writer.resetNumMaterializeCalls();
670+
writer.readChunks(dataset, datasetAttributes, ptList);
671+
672+
// one for the index, one that covers ALL the blocks)
673+
assertEquals(2, writer.getNumMaterializeCalls());
674+
675+
writer.resetNumMaterializeCalls();
676+
writer.readBlock(dataset, datasetAttributes, new long[] {0,0});
677+
// one for the index, one that covers ALL the blocks
678+
assertEquals(2, writer.getNumMaterializeCalls());
667679
}
668680
}
669681

0 commit comments

Comments
 (0)