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:
Simon Baslé
2016-05-20 14:21:47 +02:00
parent 7a57a235d8
commit 3abefe5882

View File

@@ -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);