diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/GridCacheDistributedQueryFuture.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/GridCacheDistributedQueryFuture.java index 09143f88b22a4..afd87e05766cd 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/GridCacheDistributedQueryFuture.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/GridCacheDistributedQueryFuture.java @@ -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; @@ -330,7 +332,12 @@ private void cancelPages(UUID nodeId) { } /** {@inheritDoc} */ - @Override public boolean onDone(Collection res, Throwable err) { + @Override protected boolean onDone(@Nullable Collection 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(); @@ -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; } } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/GridCacheQueryFutureAdapter.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/GridCacheQueryFutureAdapter.java index 253cf3e1b7190..bcef07fe45201 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/GridCacheQueryFutureAdapter.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/GridCacheQueryFutureAdapter.java @@ -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; } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/query/running/RunningQueryManager.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/query/running/RunningQueryManager.java index 3990d21b353ab..08be09944263a 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/query/running/RunningQueryManager.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/query/running/RunningQueryManager.java @@ -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()) { diff --git a/modules/indexing/src/test/java/org/apache/ignite/internal/processors/performancestatistics/PerformanceStatisticsQueryTest.java b/modules/indexing/src/test/java/org/apache/ignite/internal/processors/performancestatistics/PerformanceStatisticsQueryTest.java index 97661fb41b1a9..85de12e78959d 100644 --- a/modules/indexing/src/test/java/org/apache/ignite/internal/processors/performancestatistics/PerformanceStatisticsQueryTest.java +++ b/modules/indexing/src/test/java/org/apache/ignite/internal/processors/performancestatistics/PerformanceStatisticsQueryTest.java @@ -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; @@ -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);