From cf593601b97e88942162b91486907feb93d8e316 Mon Sep 17 00:00:00 2001 From: Igor Kulikov Date: Mon, 18 May 2020 16:03:51 +0300 Subject: [PATCH] Sync multi-page guava cassandra session support --- .../guava/GuavaMultiPageResultSet.java | 123 ++++++++++++++++++ .../dao/cassandra/guava/GuavaSession.java | 23 ++++ 2 files changed, 146 insertions(+) create mode 100644 common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/guava/GuavaMultiPageResultSet.java diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/guava/GuavaMultiPageResultSet.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/guava/GuavaMultiPageResultSet.java new file mode 100644 index 0000000000..dc4aebe542 --- /dev/null +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/guava/GuavaMultiPageResultSet.java @@ -0,0 +1,123 @@ +/** + * Copyright © 2016-2020 The Thingsboard Authors + * + * Licensed 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 org.thingsboard.server.dao.cassandra.guava; + +import com.datastax.oss.driver.api.core.cql.AsyncResultSet; +import com.datastax.oss.driver.api.core.cql.ColumnDefinitions; +import com.datastax.oss.driver.api.core.cql.ExecutionInfo; +import com.datastax.oss.driver.api.core.cql.ResultSet; +import com.datastax.oss.driver.api.core.cql.Row; +import com.datastax.oss.driver.api.core.cql.Statement; +import com.datastax.oss.driver.internal.core.util.CountingIterator; +import com.datastax.oss.driver.internal.core.util.concurrent.BlockingOperation; +import edu.umd.cs.findbugs.annotations.NonNull; + +import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.Iterator; +import java.util.List; + +public class GuavaMultiPageResultSet implements ResultSet { + + private final RowIterator iterator; + private final List executionInfos = new ArrayList<>(); + private ColumnDefinitions columnDefinitions; + + public GuavaMultiPageResultSet(@NonNull GuavaSession session, @NonNull Statement statement, @NonNull AsyncResultSet firstPage) { + assert firstPage.hasMorePages(); + this.iterator = new RowIterator(session, statement, firstPage); + this.executionInfos.add(firstPage.getExecutionInfo()); + this.columnDefinitions = firstPage.getColumnDefinitions(); + } + + @NonNull + @Override + public ColumnDefinitions getColumnDefinitions() { + return columnDefinitions; + } + + @NonNull + @Override + public List getExecutionInfos() { + return executionInfos; + } + + @Override + public boolean isFullyFetched() { + return iterator.isFullyFetched(); + } + + @Override + public int getAvailableWithoutFetching() { + return iterator.remaining(); + } + + @NonNull + @Override + public Iterator iterator() { + return iterator; + } + + @Override + public boolean wasApplied() { + return iterator.wasApplied(); + } + + private class RowIterator extends CountingIterator { + private GuavaSession session; + private Statement statement; + private AsyncResultSet currentPage; + private Iterator currentRows; + + private RowIterator(GuavaSession session, Statement statement, AsyncResultSet firstPage) { + super(firstPage.remaining()); + this.session = session; + this.statement = statement; + this.currentPage = firstPage; + this.currentRows = firstPage.currentPage().iterator(); + } + + @Override + protected Row computeNext() { + maybeMoveToNextPage(); + return currentRows.hasNext() ? currentRows.next() : endOfData(); + } + + private void maybeMoveToNextPage() { + if (!currentRows.hasNext() && currentPage.hasMorePages()) { + BlockingOperation.checkNotDriverThread(); + ByteBuffer nextPagingState = currentPage.getExecutionInfo().getPagingState(); + this.statement = this.statement.setPagingState(nextPagingState); + AsyncResultSet nextPage = GuavaSession.getSafe(this.session.executeAsync(this.statement)); + currentPage = nextPage; + remaining += nextPage.remaining(); + currentRows = nextPage.currentPage().iterator(); + executionInfos.add(nextPage.getExecutionInfo()); + // The definitions can change from page to page if this result set was built from a bound + // 'SELECT *', and the schema was altered. + columnDefinitions = nextPage.getColumnDefinitions(); + } + } + + private boolean isFullyFetched() { + return !currentPage.hasMorePages(); + } + + private boolean wasApplied() { + return currentPage.wasApplied(); + } + } +} diff --git a/common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/guava/GuavaSession.java b/common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/guava/GuavaSession.java index 97f300ee3c..73c8f8c80c 100644 --- a/common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/guava/GuavaSession.java +++ b/common/dao-api/src/main/java/org/thingsboard/server/dao/cassandra/guava/GuavaSession.java @@ -17,13 +17,18 @@ package org.thingsboard.server.dao.cassandra.guava; import com.datastax.oss.driver.api.core.cql.AsyncResultSet; import com.datastax.oss.driver.api.core.cql.PreparedStatement; +import com.datastax.oss.driver.api.core.cql.ResultSet; import com.datastax.oss.driver.api.core.cql.SimpleStatement; import com.datastax.oss.driver.api.core.cql.Statement; import com.datastax.oss.driver.api.core.cql.SyncCqlSession; import com.datastax.oss.driver.api.core.session.Session; import com.datastax.oss.driver.api.core.type.reflect.GenericType; import com.datastax.oss.driver.internal.core.cql.DefaultPrepareRequest; +import com.datastax.oss.driver.internal.core.cql.SinglePageResultSet; import com.google.common.util.concurrent.ListenableFuture; +import edu.umd.cs.findbugs.annotations.NonNull; + +import java.util.concurrent.ExecutionException; public interface GuavaSession extends Session, SyncCqlSession { @@ -33,6 +38,16 @@ public interface GuavaSession extends Session, SyncCqlSession { GenericType> ASYNC_PREPARED = new GenericType>() {}; + @NonNull + default ResultSet execute(@NonNull Statement statement) { + AsyncResultSet firstPage = getSafe(this.executeAsync(statement)); + if (firstPage.hasMorePages()) { + return new GuavaMultiPageResultSet(this, statement, firstPage); + } else { + return new SinglePageResultSet(firstPage); + } + } + default ListenableFuture executeAsync(Statement statement) { return this.execute(statement, ASYNC); } @@ -48,4 +63,12 @@ public interface GuavaSession extends Session, SyncCqlSession { default ListenableFuture prepareAsync(String statement) { return this.prepareAsync(SimpleStatement.newInstance(statement)); } + + static AsyncResultSet getSafe(ListenableFuture future) { + try { + return future.get(); + } catch (InterruptedException | ExecutionException e) { + throw new IllegalStateException(e); + } + } }