DATACOUCH-484 - Thread safety issue using findBy.
Support for positional parameters necessitated caching them, however that was not done in a threadsafe fashion. ThreadLocal storage is sufficient to take care of the issue. Added test for it while at it.
This commit is contained in:
@@ -50,16 +50,21 @@ import org.springframework.util.Assert;
|
|||||||
public class PartTreeN1qlBasedQuery extends AbstractN1qlBasedQuery {
|
public class PartTreeN1qlBasedQuery extends AbstractN1qlBasedQuery {
|
||||||
|
|
||||||
private final PartTree partTree;
|
private final PartTree partTree;
|
||||||
private JsonValue placeHolderValues;
|
private ThreadLocal<JsonValue> placeHolderValues;
|
||||||
|
|
||||||
public PartTreeN1qlBasedQuery(CouchbaseQueryMethod queryMethod, CouchbaseOperations couchbaseOperations) {
|
public PartTreeN1qlBasedQuery(CouchbaseQueryMethod queryMethod, CouchbaseOperations couchbaseOperations) {
|
||||||
super(queryMethod, couchbaseOperations);
|
super(queryMethod, couchbaseOperations);
|
||||||
this.partTree = new PartTree(queryMethod.getName(), queryMethod.getEntityInformation().getJavaType());
|
this.partTree = new PartTree(queryMethod.getName(), queryMethod.getEntityInformation().getJavaType());
|
||||||
|
this.placeHolderValues = new ThreadLocal<JsonValue>() {
|
||||||
|
@Override public JsonValue initialValue() {
|
||||||
|
return JsonArray.create();
|
||||||
|
}
|
||||||
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
protected JsonValue getPlaceholderValues(ParameterAccessor accessor) {
|
protected JsonValue getPlaceholderValues(ParameterAccessor accessor) {
|
||||||
return this.placeHolderValues;
|
return this.placeHolderValues.get();
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
@@ -70,7 +75,7 @@ public class PartTreeN1qlBasedQuery extends AbstractN1qlBasedQuery {
|
|||||||
N1qlCountQueryCreator queryCountCreator = new N1qlCountQueryCreator(partTree, accessor, countFrom,
|
N1qlCountQueryCreator queryCountCreator = new N1qlCountQueryCreator(partTree, accessor, countFrom,
|
||||||
getCouchbaseOperations().getConverter(), getQueryMethod());
|
getCouchbaseOperations().getConverter(), getQueryMethod());
|
||||||
Statement statement = queryCountCreator.createQuery();
|
Statement statement = queryCountCreator.createQuery();
|
||||||
this.placeHolderValues = queryCountCreator.getPlaceHolderValues();
|
this.placeHolderValues.set(queryCountCreator.getPlaceHolderValues());
|
||||||
return statement;
|
return statement;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -83,7 +88,7 @@ public class PartTreeN1qlBasedQuery extends AbstractN1qlBasedQuery {
|
|||||||
DeleteUsePath deleteUsePath = deleteFrom(bucket);
|
DeleteUsePath deleteUsePath = deleteFrom(bucket);
|
||||||
N1qlMutateQueryCreator mutateQueryCreator = new N1qlMutateQueryCreator(partTree, accessor, deleteUsePath, getCouchbaseOperations().getConverter(), getQueryMethod());
|
N1qlMutateQueryCreator mutateQueryCreator = new N1qlMutateQueryCreator(partTree, accessor, deleteUsePath, getCouchbaseOperations().getConverter(), getQueryMethod());
|
||||||
MutateLimitPath mutateFromWhereOrderBy = mutateQueryCreator.createQuery();
|
MutateLimitPath mutateFromWhereOrderBy = mutateQueryCreator.createQuery();
|
||||||
this.placeHolderValues = mutateQueryCreator.getPlaceHolderValues();
|
this.placeHolderValues.set(mutateQueryCreator.getPlaceHolderValues());
|
||||||
|
|
||||||
if (partTree.isLimiting()) {
|
if (partTree.isLimiting()) {
|
||||||
return mutateFromWhereOrderBy.limit(partTree.getMaxResults());
|
return mutateFromWhereOrderBy.limit(partTree.getMaxResults());
|
||||||
@@ -101,7 +106,7 @@ public class PartTreeN1qlBasedQuery extends AbstractN1qlBasedQuery {
|
|||||||
N1qlQueryCreator queryCreator = new N1qlQueryCreator(partTree, accessor, selectFrom,
|
N1qlQueryCreator queryCreator = new N1qlQueryCreator(partTree, accessor, selectFrom,
|
||||||
getCouchbaseOperations().getConverter(), getQueryMethod());
|
getCouchbaseOperations().getConverter(), getQueryMethod());
|
||||||
LimitPath selectFromWhereOrderBy = queryCreator.createQuery();
|
LimitPath selectFromWhereOrderBy = queryCreator.createQuery();
|
||||||
this.placeHolderValues = queryCreator.getPlaceHolderValues();
|
this.placeHolderValues.set(queryCreator.getPlaceHolderValues());
|
||||||
|
|
||||||
if (queryMethod.isPageQuery()) {
|
if (queryMethod.isPageQuery()) {
|
||||||
Pageable pageable = accessor.getPageable();
|
Pageable pageable = accessor.getPageable();
|
||||||
|
|||||||
@@ -26,7 +26,6 @@ import org.junit.runner.RunWith;
|
|||||||
|
|
||||||
import org.springframework.beans.factory.annotation.Autowired;
|
import org.springframework.beans.factory.annotation.Autowired;
|
||||||
import org.springframework.dao.DataRetrievalFailureException;
|
import org.springframework.dao.DataRetrievalFailureException;
|
||||||
;
|
|
||||||
import org.springframework.data.couchbase.ContainerResourceRunner;
|
import org.springframework.data.couchbase.ContainerResourceRunner;
|
||||||
import org.springframework.data.couchbase.IntegrationTestApplicationConfig;
|
import org.springframework.data.couchbase.IntegrationTestApplicationConfig;
|
||||||
import org.springframework.data.couchbase.repository.config.RepositoryOperationsMapping;
|
import org.springframework.data.couchbase.repository.config.RepositoryOperationsMapping;
|
||||||
@@ -41,9 +40,15 @@ import org.springframework.data.repository.core.support.RepositoryFactorySupport
|
|||||||
import org.springframework.test.context.ContextConfiguration;
|
import org.springframework.test.context.ContextConfiguration;
|
||||||
import org.springframework.test.context.TestExecutionListeners;
|
import org.springframework.test.context.TestExecutionListeners;
|
||||||
|
|
||||||
|
import java.util.ArrayList;
|
||||||
import java.util.Calendar;
|
import java.util.Calendar;
|
||||||
import java.util.Date;
|
import java.util.Date;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
|
import java.util.concurrent.Callable;
|
||||||
|
import java.util.concurrent.ExecutionException;
|
||||||
|
import java.util.concurrent.ExecutorService;
|
||||||
|
import java.util.concurrent.Executors;
|
||||||
|
import java.util.concurrent.Future;
|
||||||
import java.util.concurrent.TimeUnit;
|
import java.util.concurrent.TimeUnit;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -90,6 +95,40 @@ public class N1qlCouchbaseRepositoryIntegrationTests {
|
|||||||
try { partyRepository.deleteById(KEY_PARTY); } catch (DataRetrievalFailureException e) {}
|
try { partyRepository.deleteById(KEY_PARTY); } catch (DataRetrievalFailureException e) {}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void shouldBeThreadsafe() {
|
||||||
|
// This doesn't guarantee it, but we should catch most thread issues without
|
||||||
|
// taking too long here...
|
||||||
|
int runs = 50;
|
||||||
|
for (int i=0; i<runs; i++) {
|
||||||
|
doShouldBeThreadsafe();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
public void doShouldBeThreadsafe() {
|
||||||
|
int threads = 50;
|
||||||
|
ExecutorService service = Executors.newFixedThreadPool(threads);
|
||||||
|
List<Callable<Boolean>> callables = new ArrayList<>();
|
||||||
|
for (int thread = 0; thread < threads; ++thread) {
|
||||||
|
final int counter = thread;
|
||||||
|
Callable<Boolean> booleanSupplier = () -> {
|
||||||
|
String expectedName = "party like it's 199" + counter%12;
|
||||||
|
String foundName = partyRepository.findByName(expectedName).get(0).getName();
|
||||||
|
return expectedName.equals(foundName); //should never get false
|
||||||
|
};
|
||||||
|
callables.add(booleanSupplier);
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
List<Future<Boolean>> futures = service.invokeAll(callables);
|
||||||
|
service.shutdown();
|
||||||
|
service.awaitTermination(5, TimeUnit.SECONDS);
|
||||||
|
for (Future<Boolean> future: futures) {
|
||||||
|
assertTrue(future.get());
|
||||||
|
}
|
||||||
|
} catch (InterruptedException | ExecutionException e) {
|
||||||
|
fail("Threads failed to run " + e.getMessage());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
public void shouldFindAllWithSort() {
|
public void shouldFindAllWithSort() {
|
||||||
Iterable<Party> allByAttendanceDesc = repository.findAll(Sort.by(Sort.Direction.DESC, "attendees"));
|
Iterable<Party> allByAttendanceDesc = repository.findAll(Sort.by(Sort.Direction.DESC, "attendees"));
|
||||||
|
|||||||
@@ -42,6 +42,8 @@ public interface PartyRepository extends CouchbaseRepository<Party, String> {
|
|||||||
|
|
||||||
List<Party> findByAttendeesGreaterThanEqual(int minAttendees);
|
List<Party> findByAttendeesGreaterThanEqual(int minAttendees);
|
||||||
|
|
||||||
|
List<Party> findByName(String name);
|
||||||
|
|
||||||
List<Party> findByEventDateIs(Date targetDate);
|
List<Party> findByEventDateIs(Date targetDate);
|
||||||
|
|
||||||
@View(designDocument = "party", viewName = "byDate")
|
@View(designDocument = "party", viewName = "byDate")
|
||||||
|
|||||||
Reference in New Issue
Block a user