Skip to content

Commit b2badd1

Browse files
committed
refactor: move .stream() and .iterator() begin paginate namespace
This lets us keep all configurations 'on the left' of the operator and all operations on the right
1 parent eb7d5a2 commit b2badd1

4 files changed

Lines changed: 130 additions & 30 deletions

File tree

src/it/java/io/weaviate/integration/PaginationITest.java

Lines changed: 38 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -19,10 +19,10 @@ public class PaginationITest extends ConcurrentTest {
1919
private static WeaviateClient client = Container.WEAVIATE.getClient();
2020

2121
@Test
22-
public void test_stream() throws IOException {
22+
public void testIterateAll() throws IOException {
2323
// Arrange
2424
var nsThings = ns("Things");
25-
var count = 10;
25+
var count = 150;
2626

2727
client.collections.create(nsThings);
2828
var things = client.collections.use(nsThings);
@@ -34,8 +34,10 @@ public void test_stream() throws IOException {
3434
}
3535
assumeTrue("all objects were inserted", inserted.size() == count);
3636

37+
var allThings = things.paginate();
38+
3739
// Act: stream
38-
var gotStream = things.stream()
40+
var gotStream = allThings.stream()
3941
.map(WeaviateObject::metadata).map(WeaviateMetadata::uuid).toList();
4042

4143
// Assert
@@ -44,20 +46,47 @@ public void test_stream() throws IOException {
4446
.hasSize(inserted.size())
4547
.containsAll(inserted);
4648

47-
// Act: list
48-
var gotList = new ArrayList<String>();
49-
for (var object : things.list()) {
50-
gotList.add(object.metadata().uuid());
49+
// Act: for-loop
50+
var gotLoop = new ArrayList<String>();
51+
for (var thing : allThings) {
52+
gotLoop.add(thing.metadata().uuid());
5153
}
5254

5355
// Assert
54-
Assertions.assertThat(gotList)
56+
Assertions.assertThat(gotLoop)
5557
.as("list fetched all objects")
5658
.hasSize(inserted.size())
5759
.containsAll(inserted);
5860

5961
Assertions.assertThat(gotStream)
6062
.as("stream and list return consistent order")
61-
.containsExactlyElementsOf(gotList);
63+
.containsExactlyElementsOf(gotLoop);
64+
}
65+
66+
@Test
67+
public void testResumePagination() throws IOException {
68+
// Arrange
69+
var nsThings = ns("Things");
70+
var count = 10;
71+
72+
client.collections.create(nsThings);
73+
74+
var things = client.collections.use(nsThings);
75+
var inserted = new ArrayList<String>();
76+
for (var i = 0; i < count; i++) {
77+
var object = things.data.insert(Collections.emptyMap());
78+
inserted.add(object.metadata().uuid());
79+
}
80+
81+
// Iterate over first 5 objects
82+
String lastId = things.paginate(p -> p.pageSize(5)).stream()
83+
.limit(5).map(thing -> thing.metadata().uuid())
84+
.reduce((prev, next) -> next).get();
85+
86+
// Act
87+
var remaining = things.paginate(p -> p.resumeFrom(lastId)).stream().count();
88+
89+
// Assert
90+
Assertions.assertThat(remaining).isEqualTo(5);
6291
}
6392
}
Lines changed: 7 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -1,15 +1,13 @@
11
package io.weaviate.client6.v1.api.collections;
22

3-
import java.util.Spliterator;
4-
import java.util.Spliterators;
5-
import java.util.stream.Stream;
6-
import java.util.stream.StreamSupport;
3+
import java.util.function.Function;
74

85
import io.weaviate.client6.v1.api.collections.aggregate.WeaviateAggregateClient;
96
import io.weaviate.client6.v1.api.collections.config.WeaviateConfigClient;
107
import io.weaviate.client6.v1.api.collections.data.WeaviateDataClient;
11-
import io.weaviate.client6.v1.api.collections.query.QueryMetadata;
8+
import io.weaviate.client6.v1.api.collections.pagination.Paginator;
129
import io.weaviate.client6.v1.api.collections.query.WeaviateQueryClient;
10+
import io.weaviate.client6.v1.internal.ObjectBuilder;
1311
import io.weaviate.client6.v1.internal.grpc.GrpcTransport;
1412
import io.weaviate.client6.v1.internal.orm.CollectionDescriptor;
1513
import io.weaviate.client6.v1.internal.rest.RestTransport;
@@ -31,17 +29,11 @@ public CollectionHandle(
3129
this.aggregate = new WeaviateAggregateClient(collectionDescriptor, grpcTransport);
3230
}
3331

34-
public Stream<WeaviateObject<T, Object, QueryMetadata>> stream() {
35-
return StreamSupport.stream(spliterator(2), false);
32+
public Paginator<T> paginate() {
33+
return Paginator.of(this.query);
3634
}
3735

38-
public Iterable<WeaviateObject<T, Object, QueryMetadata>> list() {
39-
return () -> Spliterators.iterator(spliterator(2));
40-
}
41-
42-
private Spliterator<WeaviateObject<T, Object, QueryMetadata>> spliterator(int batchSize) {
43-
return new CursorSpliterator<>(batchSize,
44-
(after, limit) -> this.query.fetchObjects(
45-
query -> query.after(after).limit(limit)).objects());
36+
public Paginator<T> paginate(Function<Paginator.Builder<T>, ObjectBuilder<Paginator<T>>> fn) {
37+
return Paginator.of(this.query, fn);
4638
}
4739
}

src/main/java/io/weaviate/client6/v1/api/collections/CursorSpliterator.java renamed to src/main/java/io/weaviate/client6/v1/api/collections/pagination/CursorSpliterator.java

Lines changed: 8 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
package io.weaviate.client6.v1.api.collections;
1+
package io.weaviate.client6.v1.api.collections.pagination;
22

33
import java.util.Collections;
44
import java.util.Iterator;
@@ -7,20 +7,22 @@
77
import java.util.function.BiFunction;
88
import java.util.function.Consumer;
99

10+
import io.weaviate.client6.v1.api.collections.WeaviateObject;
1011
import io.weaviate.client6.v1.api.collections.query.QueryMetadata;
1112

12-
class CursorSpliterator<T> implements Spliterator<WeaviateObject<T, Object, QueryMetadata>> {
13-
private final int batchSize;
13+
public class CursorSpliterator<T> implements Spliterator<WeaviateObject<T, Object, QueryMetadata>> {
14+
private final int pageSize;
1415
private final BiFunction<String, Integer, List<WeaviateObject<T, Object, QueryMetadata>>> fetch;
1516

1617
// Spliterators do not promise thread-safety, so there's no mechanism
1718
// to protect access to its internal state.
1819
private String cursor;
1920
private Iterator<WeaviateObject<T, Object, QueryMetadata>> currentPage = Collections.emptyIterator();
2021

21-
public CursorSpliterator(int batchSize,
22+
public CursorSpliterator(String cursor, int pageSize,
2223
BiFunction<String, Integer, List<WeaviateObject<T, Object, QueryMetadata>>> fetch) {
23-
this.batchSize = batchSize;
24+
this.cursor = cursor;
25+
this.pageSize = pageSize;
2426
this.fetch = fetch;
2527
}
2628

@@ -33,7 +35,7 @@ public boolean tryAdvance(Consumer<? super WeaviateObject<T, Object, QueryMetada
3335
}
3436

3537
// It's OK for the cursor to be null, because it's String (object).
36-
var nextPage = fetch.apply(cursor, batchSize);
38+
var nextPage = fetch.apply(cursor, pageSize);
3739
if (nextPage.isEmpty()) {
3840
return false;
3941
}
Lines changed: 77 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,77 @@
1+
package io.weaviate.client6.v1.api.collections.pagination;
2+
3+
import java.util.Iterator;
4+
import java.util.Spliterator;
5+
import java.util.Spliterators;
6+
import java.util.function.Function;
7+
import java.util.stream.Stream;
8+
import java.util.stream.StreamSupport;
9+
10+
import io.weaviate.client6.v1.api.collections.WeaviateObject;
11+
import io.weaviate.client6.v1.api.collections.query.QueryMetadata;
12+
import io.weaviate.client6.v1.api.collections.query.WeaviateQueryClient;
13+
import io.weaviate.client6.v1.internal.ObjectBuilder;
14+
15+
public class Paginator<T> implements Iterable<WeaviateObject<T, Object, QueryMetadata>> {
16+
private static final int DEFAULT_PAGE_SIZE = 100;
17+
18+
private final WeaviateQueryClient<T> query;
19+
private final int pageSize;
20+
private final String cursor;
21+
22+
@Override
23+
public Iterator<WeaviateObject<T, Object, QueryMetadata>> iterator() {
24+
return Spliterators.iterator(spliterator());
25+
}
26+
27+
public Stream<WeaviateObject<T, Object, QueryMetadata>> stream() {
28+
return StreamSupport.stream(spliterator(), false);
29+
}
30+
31+
public Spliterator<WeaviateObject<T, Object, QueryMetadata>> spliterator() {
32+
return new CursorSpliterator<T>(cursor, pageSize,
33+
(after, limit) -> query.fetchObjects(q -> q.after(after).limit(limit)).objects());
34+
}
35+
36+
public static <T> Paginator<T> of(WeaviateQueryClient<T> query) {
37+
return of(query, ObjectBuilder.identity());
38+
}
39+
40+
public static <T> Paginator<T> of(WeaviateQueryClient<T> query,
41+
Function<Builder<T>, ObjectBuilder<Paginator<T>>> fn) {
42+
return fn.apply(new Builder<>(query)).build();
43+
}
44+
45+
Paginator(Builder<T> builder) {
46+
this.query = builder.query;
47+
this.cursor = builder.cursor;
48+
this.pageSize = builder.pageSize;
49+
}
50+
51+
public static class Builder<T> implements ObjectBuilder<Paginator<T>> {
52+
private final WeaviateQueryClient<T> query;
53+
54+
int pageSize = DEFAULT_PAGE_SIZE;
55+
String cursor;
56+
57+
public Builder(WeaviateQueryClient<T> query) {
58+
this.query = query;
59+
}
60+
61+
public Builder<T> pageSize(int pageSize) {
62+
this.pageSize = pageSize;
63+
return this;
64+
}
65+
66+
public Builder<T> resumeFrom(String uuid) {
67+
this.cursor = uuid;
68+
return this;
69+
}
70+
71+
@Override
72+
public Paginator<T> build() {
73+
return new Paginator<>(this);
74+
}
75+
}
76+
77+
}

0 commit comments

Comments
 (0)