DATACOUCH-228 - Make findByView use async client internally
This brings better performance when querying large views. Instead of looping on a synchronously acquired list of documents to produce a list of entities, the async API can be used internally to: - query the view - fetch the documents - map each document to an entity - collect into a List One can then block on that Observable to get and return the List as usual.
This commit is contained in:
@@ -40,14 +40,18 @@ import com.couchbase.client.java.query.N1qlQuery;
|
|||||||
import com.couchbase.client.java.query.N1qlQueryResult;
|
import com.couchbase.client.java.query.N1qlQueryResult;
|
||||||
import com.couchbase.client.java.query.N1qlQueryRow;
|
import com.couchbase.client.java.query.N1qlQueryRow;
|
||||||
import com.couchbase.client.java.util.features.CouchbaseFeature;
|
import com.couchbase.client.java.util.features.CouchbaseFeature;
|
||||||
|
import com.couchbase.client.java.view.AsyncViewResult;
|
||||||
|
import com.couchbase.client.java.view.AsyncViewRow;
|
||||||
import com.couchbase.client.java.view.SpatialViewQuery;
|
import com.couchbase.client.java.view.SpatialViewQuery;
|
||||||
import com.couchbase.client.java.view.SpatialViewResult;
|
import com.couchbase.client.java.view.SpatialViewResult;
|
||||||
import com.couchbase.client.java.view.SpatialViewRow;
|
import com.couchbase.client.java.view.SpatialViewRow;
|
||||||
import com.couchbase.client.java.view.ViewQuery;
|
import com.couchbase.client.java.view.ViewQuery;
|
||||||
import com.couchbase.client.java.view.ViewResult;
|
import com.couchbase.client.java.view.ViewResult;
|
||||||
import com.couchbase.client.java.view.ViewRow;
|
|
||||||
import org.slf4j.Logger;
|
import org.slf4j.Logger;
|
||||||
import org.slf4j.LoggerFactory;
|
import org.slf4j.LoggerFactory;
|
||||||
|
import rx.Observable;
|
||||||
|
import rx.functions.Func1;
|
||||||
|
|
||||||
import org.springframework.context.ApplicationEventPublisher;
|
import org.springframework.context.ApplicationEventPublisher;
|
||||||
import org.springframework.context.ApplicationEventPublisherAware;
|
import org.springframework.context.ApplicationEventPublisherAware;
|
||||||
import org.springframework.dao.OptimisticLockingFailureException;
|
import org.springframework.dao.OptimisticLockingFailureException;
|
||||||
@@ -297,36 +301,64 @@ public class CouchbaseTemplate implements CouchbaseOperations, ApplicationEventP
|
|||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public <T> List<T> findByView(ViewQuery query, Class<T> entityClass) {
|
public <T> List<T> findByView(ViewQuery query, final Class<T> entityClass) {
|
||||||
//we'll always need to get documents, as a RawJsonDocument, so we should force includeDocs(false)
|
//we'll always need to get documents, as a RawJsonDocument, so we should force that target class
|
||||||
//so that the caller doesn't set a bad target class unintentionally, pre-loading with a bad type.
|
//so that the caller doesn't set a bad target class unintentionally, pre-loading with a bad type.
|
||||||
query.includeDocs(false);
|
//TODO DATACOUCH-227 reproduce retainOrder parameter
|
||||||
|
if (!query.isIncludeDocs() || !query.includeDocsTarget().equals(RawJsonDocument.class)) {
|
||||||
|
query.includeDocs(RawJsonDocument.class);
|
||||||
|
}
|
||||||
//we'll always map the document to the entity, hence reduce never makes sense.
|
//we'll always map the document to the entity, hence reduce never makes sense.
|
||||||
query.reduce(false);
|
query.reduce(false);
|
||||||
|
|
||||||
try {
|
return executeAsync(client.async().query(query))
|
||||||
final ViewResult response = queryView(query);
|
.flatMap(new Func1<AsyncViewResult, Observable<AsyncViewRow>>() {
|
||||||
if (response.error() != null) {
|
@Override
|
||||||
throw new CouchbaseQueryExecutionException("Unable to execute view query due to the following view error: " +
|
public Observable<AsyncViewRow> call(AsyncViewResult asyncViewResult) {
|
||||||
response.error().toString());
|
return asyncViewResult
|
||||||
}
|
.error()
|
||||||
|
.flatMap(new Func1<JsonObject, Observable<AsyncViewRow>>() {
|
||||||
List<ViewRow> allRows = response.allRows();
|
@Override
|
||||||
|
public Observable<AsyncViewRow> call(JsonObject error) {
|
||||||
final List<T> result = new ArrayList<T>(allRows.size());
|
return Observable.error(new CouchbaseQueryExecutionException("Unable to execute view query due to the following view error: " + error.toString()));
|
||||||
for (final ViewRow row : allRows) {
|
}})
|
||||||
//cope with potential weak consistency and deletions
|
.switchIfEmpty(asyncViewResult.rows());
|
||||||
T entity = mapToEntity(row.id(), row.document(RawJsonDocument.class), entityClass);
|
}
|
||||||
if (entity != null) {
|
})
|
||||||
result.add(entity);
|
.flatMap(new Func1<AsyncViewRow, Observable<T>>() {
|
||||||
}
|
@Override
|
||||||
}
|
public Observable<T> call(AsyncViewRow row) {
|
||||||
|
final String id = row.id();
|
||||||
return result;
|
return row
|
||||||
}
|
.document(RawJsonDocument.class)
|
||||||
catch (TranscodingException e) {
|
.map(new Func1<RawJsonDocument, T>() {
|
||||||
throw new CouchbaseQueryExecutionException("Unable to execute view query", e);
|
@Override
|
||||||
}
|
public T call(RawJsonDocument rawJsonDocument) {
|
||||||
|
//cope with potential weak consistency and deletions
|
||||||
|
T entity = mapToEntity(id, rawJsonDocument, entityClass);
|
||||||
|
return entity;
|
||||||
|
}
|
||||||
|
});
|
||||||
|
}})
|
||||||
|
.filter(new Func1<T, Boolean>() {
|
||||||
|
@Override
|
||||||
|
public Boolean call(T t) {
|
||||||
|
return t != null;
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.onErrorResumeNext(new Func1<Throwable, Observable<T>>() {
|
||||||
|
@Override
|
||||||
|
public Observable<T> call(Throwable throwable) {
|
||||||
|
if (throwable instanceof TranscodingException) {
|
||||||
|
return Observable.error(new CouchbaseQueryExecutionException("Unable to execute view query", throwable));
|
||||||
|
} else {
|
||||||
|
return Observable.error(throwable);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.toList()
|
||||||
|
.toBlocking()
|
||||||
|
.single();
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
@@ -508,6 +540,26 @@ public class CouchbaseTemplate implements CouchbaseOperations, ApplicationEventP
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public <T> Observable<T> executeAsync(Observable<T> asyncAction) {
|
||||||
|
return asyncAction
|
||||||
|
.onErrorResumeNext(new Func1<Throwable, Observable<T>>() {
|
||||||
|
@Override
|
||||||
|
public Observable<T> call(Throwable e) {
|
||||||
|
if (e instanceof RuntimeException) {
|
||||||
|
return Observable.error(exceptionTranslator.translateExceptionIfPossible((RuntimeException) e));
|
||||||
|
} else if (e instanceof TimeoutException) {
|
||||||
|
return Observable.error(new QueryTimeoutException(e.getMessage(), e));
|
||||||
|
} else if (e instanceof InterruptedException) {
|
||||||
|
return Observable.error(new OperationInterruptedException(e.getMessage(), e));
|
||||||
|
} else if (e instanceof ExecutionException) {
|
||||||
|
return Observable.error(new OperationInterruptedException(e.getMessage(), e));
|
||||||
|
} else {
|
||||||
|
return Observable.error(e);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
private void doPersist(Object objectToPersist, final PersistTo persistTo, final ReplicateTo replicateTo,
|
private void doPersist(Object objectToPersist, final PersistTo persistTo, final ReplicateTo replicateTo,
|
||||||
final PersistType persistType) {
|
final PersistType persistType) {
|
||||||
ensureNotIterable(objectToPersist);
|
ensureNotIterable(objectToPersist);
|
||||||
|
|||||||
Reference in New Issue
Block a user