Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -37,10 +37,12 @@
import org.apache.ignite.internal.processors.cache.query.reducer.NodePageStream;
import org.apache.ignite.internal.processors.cache.query.reducer.TextQueryReducer;
import org.apache.ignite.internal.processors.cache.query.reducer.UnsortedCacheQueryReducer;
import org.apache.ignite.internal.processors.query.QueryUtils;
import org.apache.ignite.internal.thread.context.concurrent.IgniteCompletableFuture;
import org.apache.ignite.internal.util.lang.GridPlainCallable;
import org.apache.ignite.internal.util.typedef.F;
import org.apache.ignite.internal.util.typedef.internal.U;
import org.jetbrains.annotations.Nullable;

import static org.apache.ignite.internal.processors.cache.query.GridCacheQueryType.INDEX;
import static org.apache.ignite.internal.processors.cache.query.GridCacheQueryType.TEXT;
Expand Down Expand Up @@ -330,7 +332,12 @@ private void cancelPages(UUID nodeId) {
}

/** {@inheritDoc} */
@Override public boolean onDone(Collection<R> res, Throwable err) {
@Override protected boolean onDone(@Nullable Collection<R> res, @Nullable Throwable err, boolean cancel) {
boolean done = super.onDone(res, err, cancel);

if (!done)
return false;

if (cctx.kernalContext().performanceStatistics().enabled() && startTimeNanos > 0) {
GridCacheQueryType type = qry.query().type();

Expand All @@ -347,9 +354,9 @@ private void cancelPages(UUID nodeId) {
reqId,
startTimeNanos,
System.nanoTime() - startTimeNanos,
err == null);
err == null || QueryUtils.wasCancelled(err));
}

return super.onDone(res, err);
return true;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -344,7 +344,7 @@ void clear() {
/** {@inheritDoc} */
@Override public boolean cancel() throws IgniteCheckedException {
if (onCancelled()) {
cancelQuery(new IgniteCheckedException("Query was cancelled."));
cancelQuery(new QueryCancelledException());

return true;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -503,7 +503,7 @@ public void unregister(long qryId, @Nullable Throwable failReason) {
qry.id(),
qry.startTime(),
System.nanoTime() - qry.startTimeNanos(),
!failed);
!failed || QueryUtils.wasCancelled(failReason));
}

if (!qryFinishedListeners.isEmpty()) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@
import org.apache.ignite.cache.query.FieldsQueryCursor;
import org.apache.ignite.cache.query.IndexQuery;
import org.apache.ignite.cache.query.Query;
import org.apache.ignite.cache.query.QueryCursor;
import org.apache.ignite.cache.query.ScanQuery;
import org.apache.ignite.cache.query.SqlFieldsQuery;
import org.apache.ignite.client.Config;
Expand Down Expand Up @@ -264,6 +265,59 @@ public void testSqlFieldsLocalQuery() throws Exception {
assertEquals("local", flags.get());
}

/** @throws Exception If failed. */
@Test
public void testScanQueryCursorNotFullyRead() throws Exception {
checkCursorNotFullyRead(new ScanQuery<>().setPageSize(pageSize));
}

/** @throws Exception If failed. */
@Test
public void testIndexQueryCursorNotFullyRead() throws Exception {
checkCursorNotFullyRead(new IndexQuery<>(Integer.class).setPageSize(pageSize));
}

/** @throws Exception If failed. */
@Test
public void testSqlFieldsQueryCursorNotFullyRead() throws Exception {
checkCursorNotFullyRead(new SqlFieldsQuery("select * from " + DEFAULT_CACHE_NAME).setPageSize(pageSize));
}

/** Checks that a query is successful when its cursor is closed before all rows are read. */
private void checkCursorNotFullyRead(Query<?> qry) throws Exception {
Assume.assumeTrue("Query result fits into one page.", pageSize < ENTRY_COUNT);

cleanPerformanceStatisticsDir();

startCollectStatistics();

QueryCursor<?> cursor;

if (clientType == SERVER)
cursor = srv.cache(DEFAULT_CACHE_NAME).query(qry);
else if (clientType == CLIENT)
cursor = client.cache(DEFAULT_CACHE_NAME).query(qry);
else
cursor = thinClient.cache(DEFAULT_CACHE_NAME).query(qry);

cursor.iterator().next();

cursor.close();

AtomicInteger qryCnt = new AtomicInteger();

stopCollectStatisticsAndRead(new TestHandler() {
@Override public void query(UUID nodeId, GridCacheQueryType type, String text, long id, long queryStartTime,
long duration, boolean success) {
qryCnt.incrementAndGet();

assertTrue(success);
}
});

assertEquals(1, qryCnt.get());
}

/** Check query. */
private void checkQuery(GridCacheQueryType type, Query<?> qry, String text, boolean hasReducer) throws Exception {
client.cluster().state(INACTIVE);
Expand Down
Loading