Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
67 changes: 65 additions & 2 deletions examples/src/main/java/io/milvus/v2/IteratorExample.java
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
import com.google.gson.Gson;
import com.google.gson.JsonObject;
import io.milvus.orm.iterator.QueryIterator;
import io.milvus.orm.iterator.QueryIteratorCursor;
import io.milvus.orm.iterator.SearchIterator;
import io.milvus.orm.iterator.SearchIteratorV2;
import io.milvus.response.QueryResultsWrapper;
Expand Down Expand Up @@ -162,7 +163,7 @@ private static void queryIterator(String expr, int batchSize, int offset, int li
}

private static void queryIteratorWithTemplate(int batchSize) {
System.out.println("\n========== queryIterator() ==========");
System.out.println("\n========== queryIteratorWithTemplate() ==========");
List<Long> ids = new ArrayList<>();
for (long i = 500L; i < 600L; i++) {
ids.add(i);
Expand Down Expand Up @@ -199,6 +200,67 @@ private static void queryIteratorWithTemplate(int batchSize) {
}


// Query iterator with a resumable cursor
private static void queryIteratorWithCursor(int batchSize, int limit) {
System.out.println("\n========== queryIteratorWithCursor() with resumable cursor ==========");
String expr = ID_FIELD + " < 3000";

QueryIterator queryIterator = client.queryIterator(QueryIteratorReq.builder()
.collectionName(COLLECTION_NAME)
.expr(expr)
.outputFields(Lists.newArrayList(ID_FIELD, AGE_FIELD))
.batchSize(batchSize)
.limit(limit)
.consistencyLevel(ConsistencyLevel.BOUNDED)
.build());

// Read four pages, then capture the cursor and close the iterator. With a batch size
// of 10, these pages contain 40 records in total.
System.out.println("QueryIterator first four pages results:");
int returnedBeforeCursor = 0;
for (int page = 0; page < 4; page++) {
List<QueryResultsWrapper.RowRecord> records = queryIterator.next();
for (QueryResultsWrapper.RowRecord record : records) {
System.out.println(record);
returnedBeforeCursor++;
}
}

QueryIteratorCursor cursor = queryIterator.getCursor();
queryIterator.close();
System.out.printf("%d query results returned before cursor capture, cursor=%s%n",
returnedBeforeCursor, cursor);

// resume pagination from the captured cursor in a brand new iterator. The limit is
// per-iterator, so subtract the rows already consumed to keep the combined total = limit.
QueryIterator resumed = client.queryIterator(QueryIteratorReq.builder()
.collectionName(COLLECTION_NAME)
.expr(expr)
.outputFields(Lists.newArrayList(ID_FIELD, AGE_FIELD))
.batchSize(batchSize)
.limit(limit - returnedBeforeCursor)
.cursor(cursor)
.consistencyLevel(ConsistencyLevel.BOUNDED)
.build());

System.out.println("Resumed iterator results:");
int counter = 0;
while (true) {
List<QueryResultsWrapper.RowRecord> res = resumed.next();
if (res.isEmpty()) {
System.out.println("resumed query iteration finished, close");
resumed.close();
break;
}
for (QueryResultsWrapper.RowRecord record : res) {
System.out.println(record);
counter++;
}
}
System.out.printf("%d query results returned after resume%n", counter);
}


// Search iterator V1
private static void searchIteratorV1(String expr, String params, int batchSize, int topK) {
System.out.println("\n========== searchIteratorV1() ==========");
Expand Down Expand Up @@ -275,7 +337,7 @@ private static void searchIteratorV2(String filter, Map<String, Object> params,
}

private static void searchIteratorV2WithTemplate(int batchSize) {
System.out.println("\n========== searchIteratorV2() ==========");
System.out.println("\n========== searchIteratorV2WithTemplate() ==========");
List<Long> ids = new ArrayList<>();
for (long i = 500L; i < 600L; i++) {
ids.add(i);
Expand Down Expand Up @@ -323,6 +385,7 @@ public static void main(String[] args) {

queryIterator("userID < 3000", 1, 5, 10000);
queryIteratorWithTemplate(80);
queryIteratorWithCursor(10, 100);

searchIteratorV1("userAge > 50 &&userAge < 100", "{\"range_filter\": 15.0, \"radius\": 20.0}", 100, 500);
searchIteratorV1("", "", 1, 3000);
Expand Down
35 changes: 35 additions & 0 deletions sdk-core/src/main/java/io/milvus/orm/iterator/QueryIterator.java
Original file line number Diff line number Diff line change
Expand Up @@ -118,6 +118,12 @@ public QueryIterator(QueryIteratorReq queryIteratorReq,
this.vectorUtils.setEndpoint(blockingStub.getEndpoint());
this.vectorUtils.setCurrentDbName(blockingStub.getDatabaseName());

if (queryIteratorReq.getCursor() != null) {
// resume from a previously captured cursor: reuse its session ts and pk/element
// position, and skip the setup-ts query and offset seek.
restoreFromCursor(queryIteratorReq.getCursor());
return;
}
setupTsByRequest();
seek();
}
Expand Down Expand Up @@ -189,6 +195,35 @@ public void close() {
iteratorCache.releaseCache(cacheIdInUse);
}

/**
* Capture the current iterator position as a resumable cursor. The cursor holds the
* session timestamp, the last primary key returned, and (for element-filter iterators)
* the last matched element offset, so pagination can continue in a new iterator.
*/
public QueryIteratorCursor getCursor() {
QueryIteratorCursor.QueryIteratorCursorBuilder builder = QueryIteratorCursor.builder()
.sessionTs(sessionTs)
.lastElementOffset(nextElementOffset == null ? null
: ((Number) nextElementOffset).longValue());
if (nextId == null) {
return builder.build();
}
if (primaryField.getDataType() == DataType.VarChar) {
builder.strPk(String.valueOf(nextId));
} else {
builder.intPk(((Number) nextId).longValue());
}
return builder.build();
}

private void restoreFromCursor(QueryIteratorCursor cursor) {
this.sessionTs = cursor.getSessionTs();
this.nextId = cursor.getStrPk() != null ? cursor.getStrPk() : cursor.getIntPk();
Comment thread
yhmo marked this conversation as resolved.
Comment thread
yhmo marked this conversation as resolved.
this.nextElementOffset = cursor.getLastElementOffset();
this.cacheIdInUse = NO_CACHE_ID;
this.offset = 0;
}

private void updateCursor(List<QueryResultsWrapper.RowRecord> res) {
if (res.isEmpty()) {
return;
Expand Down
154 changes: 154 additions & 0 deletions sdk-core/src/main/java/io/milvus/orm/iterator/QueryIteratorCursor.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,154 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/

package io.milvus.orm.iterator;

import io.milvus.exception.ParamException;
import io.milvus.grpc.QueryCursor;

/**
* A serializable snapshot of a {@link QueryIterator} position, used to resume
* pagination from where a previous iterator left off (pymilvus parity:
* {@code QueryIterator.get_cursor()} / {@code QueryIteratorCursor}).
*
* <p>The cursor captures the session timestamp of the original iterator, the last
* primary key returned, and (for element-filter iterators) the last matched element
* offset. Resume by passing it back through {@code QueryIteratorReq.cursor(...)}.
*/
public class QueryIteratorCursor {
private final long sessionTs;
private final Long intPk;
private final String strPk;
private final Long lastElementOffset;

private QueryIteratorCursor(QueryIteratorCursorBuilder builder) {
this.sessionTs = builder.sessionTs;
this.intPk = builder.intPk;
this.strPk = builder.strPk;
this.lastElementOffset = builder.lastElementOffset;
}

public static QueryIteratorCursorBuilder builder() {
return new QueryIteratorCursorBuilder();
}

public long getSessionTs() {
return sessionTs;
}

public Long getIntPk() {
return intPk;
}

public String getStrPk() {
return strPk;
}

public Long getLastElementOffset() {
return lastElementOffset;
}

/**
* Serialize this cursor to the gRPC {@code QueryCursor} message.
*
* <p>The gRPC {@code QueryCursor} message has no field for the last element offset, so
* serializing an element-filter cursor would silently drop the element position and cause
* the resumed iterator to skip rows. Reject such cursors to avoid the data loss.
*
* @throws ParamException if this cursor carries an element offset
*/
public QueryCursor toProto() {
if (lastElementOffset != null) {
throw new ParamException("Cannot serialize an element-filter cursor to QueryCursor: "
+ "the gRPC message has no element-offset field, resuming would skip rows. "
+ "Resume in memory via QueryIteratorReq.cursor(...) instead.");
}
QueryCursor.Builder builder = QueryCursor.newBuilder().setSessionTs(sessionTs);
if (strPk != null) {
builder.setStrPk(strPk);
} else if (intPk != null) {
builder.setIntPk(intPk);
}
return builder.build();
}

/**
* Reconstruct a cursor from a gRPC {@code QueryCursor} message, or {@code null}
* when the input is null.
*/
public static QueryIteratorCursor fromProto(QueryCursor cursor) {
if (cursor == null) {
return null;
}
QueryIteratorCursorBuilder builder = QueryIteratorCursor.builder()
.sessionTs(cursor.getSessionTs());
switch (cursor.getCursorPkCase()) {
case STR_PK:
builder.strPk(cursor.getStrPk());
break;
case INT_PK:
builder.intPk(cursor.getIntPk());
break;
default:
break;
}
return builder.build();
}

@Override
public String toString() {
return "QueryIteratorCursor{" +
"sessionTs=" + sessionTs +
", intPk=" + intPk +
", strPk='" + strPk + '\'' +
", lastElementOffset=" + lastElementOffset +
'}';
}

public static class QueryIteratorCursorBuilder {
private long sessionTs;
private Long intPk;
private String strPk;
private Long lastElementOffset;

public QueryIteratorCursorBuilder sessionTs(long sessionTs) {
this.sessionTs = sessionTs;
return this;
}

public QueryIteratorCursorBuilder intPk(Long intPk) {
this.intPk = intPk;
return this;
Comment thread
yhmo marked this conversation as resolved.
}

public QueryIteratorCursorBuilder strPk(String strPk) {
this.strPk = strPk;
return this;
}

public QueryIteratorCursorBuilder lastElementOffset(Long lastElementOffset) {
this.lastElementOffset = lastElementOffset;
return this;
}

public QueryIteratorCursor build() {
return new QueryIteratorCursor(this);
}
}
}
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package io.milvus.v2.service.vector.request;

import com.google.common.collect.Lists;
import io.milvus.orm.iterator.QueryIteratorCursor;
import io.milvus.v2.common.ConsistencyLevel;

import java.util.HashMap;
Expand Down Expand Up @@ -37,6 +38,11 @@ public class QueryIteratorReq {
// Boolean, Long, Double, String, List<Boolean>, List<Long>, List<Double>, List<String>
private Map<String, Object> filterTemplateValues;

// A previously captured cursor to resume pagination from (pymilvus parity).
// When set, the iterator continues from the cursor's session ts and pk/element
// position instead of starting over; offset is ignored in that case.
private QueryIteratorCursor cursor;

private QueryIteratorReq(QueryIteratorReqBuilder builder) {
this.databaseName = builder.databaseName;
this.collectionName = builder.collectionName;
Expand All @@ -52,6 +58,7 @@ private QueryIteratorReq(QueryIteratorReqBuilder builder) {
this.batchSize = builder.batchSize;
this.reduceStopForBest = builder.reduceStopForBest;
this.filterTemplateValues = builder.filterTemplateValues;
this.cursor = builder.cursor;
}

public static QueryIteratorReqBuilder builder() {
Expand Down Expand Up @@ -172,6 +179,14 @@ public Map<String, Object> getFilterTemplateValues() {
return filterTemplateValues;
}

public QueryIteratorCursor getCursor() {
return cursor;
Comment thread
yhmo marked this conversation as resolved.
}

public void setCursor(QueryIteratorCursor cursor) {
this.cursor = cursor;
}

@Override
public String toString() {
return "QueryIteratorReq{" +
Expand All @@ -188,6 +203,7 @@ public String toString() {
", timezone='" + timezone + '\'' +
", batchSize=" + batchSize +
", reduceStopForBest=" + reduceStopForBest +
", cursor=" + cursor +
'}';
}

Expand All @@ -206,6 +222,7 @@ public static class QueryIteratorReqBuilder {
private long batchSize = 1000L;
private boolean reduceStopForBest = true;
private Map<String, Object> filterTemplateValues = new HashMap<>();
private QueryIteratorCursor cursor;

public QueryIteratorReqBuilder databaseName(String databaseName) {
this.databaseName = databaseName;
Expand Down Expand Up @@ -282,6 +299,11 @@ public QueryIteratorReqBuilder filterTemplateValues(Map<String, Object> filterTe
return this;
}

public QueryIteratorReqBuilder cursor(QueryIteratorCursor cursor) {
this.cursor = cursor;
return this;
}

public QueryIteratorReq build() {
return new QueryIteratorReq(this);
}
Expand Down
4 changes: 2 additions & 2 deletions sdk-core/src/test/java/io/milvus/docker-compose-multi.yml
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@ services:

standalone:
container_name: milvus-javasdk-standalone-1
image: milvusdb/milvus:master-20260723-32a44262
image: milvusdb/milvus:master-20260825-dd1ee671
user: "0:0"
command: ["milvus", "run", "standalone"]
security_opt:
Expand Down Expand Up @@ -83,7 +83,7 @@ services:

standaloneslave:
container_name: milvus-javasdk-standalone-2
image: milvusdb/milvus:master-20260723-32a44262
image: milvusdb/milvus:master-20260825-dd1ee671
user: "0:0"
command: ["milvus", "run", "standalone"]
security_opt:
Expand Down
Loading
Loading