From 811c024041cbcd380ca9d675ce4d39c16c4bc5d9 Mon Sep 17 00:00:00 2001 From: Anastasiia Smirnova Date: Tue, 10 May 2016 14:27:31 +0300 Subject: [PATCH] DATACOUCH-224 - Improve concurrent save/insert, optimistic locking --- .../data/couchbase/core/AsyncUtils.java | 20 ++++ .../core/CouchbaseTemplateTests.java | 96 ++++++++++++++++++- .../SimpleCouchbaseRepositoryTests.java | 93 ++++++++++++++++-- .../couchbase/core/CouchbaseTemplate.java | 19 +++- 4 files changed, 214 insertions(+), 14 deletions(-) create mode 100644 src/integration/java/org/springframework/data/couchbase/core/AsyncUtils.java diff --git a/src/integration/java/org/springframework/data/couchbase/core/AsyncUtils.java b/src/integration/java/org/springframework/data/couchbase/core/AsyncUtils.java new file mode 100644 index 00000000..8df3f22e --- /dev/null +++ b/src/integration/java/org/springframework/data/couchbase/core/AsyncUtils.java @@ -0,0 +1,20 @@ +package org.springframework.data.couchbase.core; + +import java.util.Collection; +import java.util.Collections; +import java.util.List; +import java.util.concurrent.*; + +public class AsyncUtils { + + public static void executeConcurrently(int numThreads, Callable task) throws Exception { + ExecutorService pool = Executors.newFixedThreadPool(numThreads); + + Collection> tasks = Collections.nCopies(numThreads, task); + + List> futures = pool.invokeAll(tasks); + for (Future future : futures) { + future.get(numThreads, TimeUnit.SECONDS); + } + } +} diff --git a/src/integration/java/org/springframework/data/couchbase/core/CouchbaseTemplateTests.java b/src/integration/java/org/springframework/data/couchbase/core/CouchbaseTemplateTests.java index 7a0222c5..9668a851 100644 --- a/src/integration/java/org/springframework/data/couchbase/core/CouchbaseTemplateTests.java +++ b/src/integration/java/org/springframework/data/couchbase/core/CouchbaseTemplateTests.java @@ -26,12 +26,20 @@ import static org.junit.Assert.*; import java.util.ArrayList; import java.util.Arrays; +import java.util.Collection; +import java.util.Collections; import java.util.Date; import java.util.HashMap; import java.util.LinkedList; import java.util.List; import java.util.Map; import java.util.Random; +import java.util.concurrent.Callable; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; import com.couchbase.client.java.Bucket; import com.couchbase.client.java.document.RawJsonDocument; @@ -40,12 +48,13 @@ import com.couchbase.client.java.query.N1qlParams; import com.couchbase.client.java.query.N1qlQuery; import com.couchbase.client.java.query.N1qlQueryResult; import com.couchbase.client.java.query.consistency.ScanConsistency; -import com.couchbase.client.java.query.dsl.Expression; import com.couchbase.client.java.view.Stale; import com.couchbase.client.java.view.ViewQuery; import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.ObjectMapper; +import org.junit.Rule; import org.junit.Test; +import org.junit.rules.TestName; import org.junit.runner.RunWith; import org.springframework.beans.factory.annotation.Autowired; @@ -69,6 +78,9 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @TestExecutionListeners(CouchbaseTemplateViewListener.class) public class CouchbaseTemplateTests { + @Rule + public TestName testName = new TestName(); + @Autowired private Bucket client; @@ -132,7 +144,11 @@ public class CouchbaseTemplateTests { assertEquals("Mr. A", resultConv.get("name")); doc = new SimplePerson(id, "Mr. B"); - template.insert(doc); + try { + template.insert(doc); + } catch (OptimisticLockingFailureException e) { + //ignore, since this insert should fail + } resultDoc = client.get(id, RawJsonDocument.class); assertNotNull(resultDoc); @@ -315,7 +331,11 @@ public class CouchbaseTemplateTests { VersionedClass versionedClass = new VersionedClass("versionedClass:2", "foobar"); template.insert(versionedClass); long version1 = versionedClass.getVersion(); - template.insert(versionedClass); + try { + template.insert(versionedClass); + } catch (OptimisticLockingFailureException e) { + //ignore, since this insert should fail + } long version2 = versionedClass.getVersion(); assertTrue(version1 > 0); @@ -402,6 +422,67 @@ public class CouchbaseTemplateTests { assertEquals(versionedClass.getVersion(), foundClass.getVersion()); } + @Test + public void shouldUpdateAlreadyExistingDocument() throws Exception { + final String key = testName.getMethodName(); + removeIfExist(key); + + final AtomicLong counter = new AtomicLong(); + + VersionedClass initial = new VersionedClass(key, "value-0"); + template.save(initial); + + AsyncUtils.executeConcurrently(3, new Callable() { + @Override + public Void call() throws Exception { + boolean saved = false; + while(!saved) { + long counterValue = counter.incrementAndGet(); + VersionedClass messageData = template.findById(key, VersionedClass.class); + messageData.field = "value-" + counterValue; + try { + template.save(messageData); + saved = true; + } catch (OptimisticLockingFailureException e) { + } + } + return null; + } + }); + + VersionedClass actual = template.findById(key, VersionedClass.class); + + assertNotEquals(initial.field, actual.field); + assertNotEquals(initial.version, actual.version); + } + + @Test + public void shouldInsertOnlyFirstDocumentAndNextAttemptsShouldFailWithOptimisticLockingException() throws Exception { + final String key = testName.getMethodName(); + removeIfExist(key); + + final AtomicLong counter = new AtomicLong(); + final AtomicLong optimisticLockCounter = new AtomicLong(); + AsyncUtils.executeConcurrently(5, new Callable() { + @Override + public Void call() throws Exception { + long counterValue = counter.incrementAndGet(); + String data = "value-" + counterValue; + VersionedClass messageData = new VersionedClass(key, data); + try { + template.insert(messageData); + } catch (OptimisticLockingFailureException e) { + optimisticLockCounter.incrementAndGet(); + } + //should save operation throw OptimisticLockingFailureException on next attempts to save? + return null; + } + }); + + + assertEquals(4, optimisticLockCounter.intValue()); + } + /** * @see DATACOUCH-59 */ @@ -647,6 +728,15 @@ public class CouchbaseTemplateTests { public void setField(String field) { this.field = field; } + + @Override + public String toString() { + return "VersionedClass{" + + "id='" + id + '\'' + + ", version=" + version + + ", field='" + field + '\'' + + '}'; + } } @Document diff --git a/src/integration/java/org/springframework/data/couchbase/repository/SimpleCouchbaseRepositoryTests.java b/src/integration/java/org/springframework/data/couchbase/repository/SimpleCouchbaseRepositoryTests.java index 9dbf8687..b333cde2 100644 --- a/src/integration/java/org/springframework/data/couchbase/repository/SimpleCouchbaseRepositoryTests.java +++ b/src/integration/java/org/springframework/data/couchbase/repository/SimpleCouchbaseRepositoryTests.java @@ -16,27 +16,26 @@ package org.springframework.data.couchbase.repository; -import static org.junit.Assert.*; - -import java.util.Arrays; -import java.util.List; - import com.couchbase.client.java.Bucket; import com.couchbase.client.java.document.JsonDocument; import com.couchbase.client.java.error.CASMismatchException; +import com.couchbase.client.java.error.DocumentDoesNotExistException; import com.couchbase.client.java.view.Stale; import com.couchbase.client.java.view.ViewQuery; import org.junit.Before; import org.junit.Ignore; +import org.junit.Rule; import org.junit.Test; +import org.junit.rules.TestName; import org.junit.runner.RunWith; - import org.springframework.beans.factory.annotation.Autowired; import org.springframework.dao.OptimisticLockingFailureException; import org.springframework.data.annotation.Id; import org.springframework.data.annotation.Version; import org.springframework.data.couchbase.IntegrationTestApplicationConfig; +import org.springframework.data.couchbase.core.AsyncUtils; import org.springframework.data.couchbase.core.CouchbaseQueryExecutionException; +import org.springframework.data.couchbase.core.CouchbaseTemplateTests; import org.springframework.data.couchbase.core.mapping.Document; import org.springframework.data.couchbase.repository.config.RepositoryOperationsMapping; import org.springframework.data.couchbase.repository.support.CouchbaseRepositoryFactory; @@ -46,6 +45,15 @@ import org.springframework.test.context.ContextConfiguration; import org.springframework.test.context.TestExecutionListeners; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import java.util.Arrays; +import java.util.Collection; +import java.util.Collections; +import java.util.List; +import java.util.concurrent.*; +import java.util.concurrent.atomic.AtomicLong; + +import static org.junit.Assert.*; + /** * @author Michael Nitschinger */ @@ -54,6 +62,9 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @TestExecutionListeners(SimpleCouchbaseRepositoryListener.class) public class SimpleCouchbaseRepositoryTests { + @Rule + public TestName testName = new TestName(); + @Autowired private Bucket client; @@ -73,6 +84,13 @@ public class SimpleCouchbaseRepositoryTests { versionedDataRepository = factory.getRepository(VersionedDataRepository.class); } + private void remove(String key) { + try { + client.remove(key); + } catch (DocumentDoesNotExistException e) { + } + } + @Test public void simpleCrud() { String key = "my_unique_user_key"; @@ -220,6 +238,69 @@ public class SimpleCouchbaseRepositoryTests { } } + @Test + public void shouldUpdateDocumentConcurrently() throws Exception { + final String key = testName.getMethodName(); + remove(key); + + final AtomicLong counter = new AtomicLong(); + final AtomicLong updatedCounter = new AtomicLong(); + VersionedData initial = new VersionedData(key, "value-initial"); + versionedDataRepository.save(initial); + assertNotEquals(0L, initial.version); + + Callable task = new Callable() { + @Override + public Void call() throws Exception { + boolean updated = false; + while(!updated) { + long counterValue = counter.incrementAndGet(); + VersionedData messageData = versionedDataRepository.findOne(key); + messageData.data = "value-" + counterValue; + try { + versionedDataRepository.save(messageData); + updated = true; + updatedCounter.incrementAndGet(); + } catch (OptimisticLockingFailureException e) { + } + } + return null; + } + }; + AsyncUtils.executeConcurrently(5, task); + + assertNotEquals(initial.data, versionedDataRepository.findOne(key).data); + assertEquals(5, updatedCounter.intValue()); + } + + @Test + public void shouldFailOnMultipleConcurrentSaves() throws Exception { + final String key = testName.getMethodName(); + remove(key); + + final AtomicLong counter = new AtomicLong(); + final AtomicLong optimisticLockCounter = new AtomicLong(); + + Callable task = new Callable() { + @Override + public Void call() throws Exception { + long counterValue = counter.incrementAndGet(); + VersionedData messageData = new VersionedData(key, "value-" + counterValue); + try { + versionedDataRepository.save(messageData); + } catch (OptimisticLockingFailureException e) { + optimisticLockCounter.incrementAndGet(); + } + return null; + } + }; + + AsyncUtils.executeConcurrently(5, task); + + assertEquals(4, optimisticLockCounter.intValue()); + } + + public interface VersionedDataRepository extends CouchbaseRepository { } @Document diff --git a/src/main/java/org/springframework/data/couchbase/core/CouchbaseTemplate.java b/src/main/java/org/springframework/data/couchbase/core/CouchbaseTemplate.java index f97105f7..cc89f90f 100644 --- a/src/main/java/org/springframework/data/couchbase/core/CouchbaseTemplate.java +++ b/src/main/java/org/springframework/data/couchbase/core/CouchbaseTemplate.java @@ -35,7 +35,6 @@ import com.couchbase.client.java.document.RawJsonDocument; import com.couchbase.client.java.document.json.JsonObject; import com.couchbase.client.java.error.CASMismatchException; import com.couchbase.client.java.error.DocumentAlreadyExistsException; -import com.couchbase.client.java.error.DocumentDoesNotExistException; import com.couchbase.client.java.error.TranscodingException; import com.couchbase.client.java.query.N1qlQuery; import com.couchbase.client.java.query.N1qlQueryResult; @@ -49,7 +48,6 @@ import com.couchbase.client.java.view.ViewResult; import com.couchbase.client.java.view.ViewRow; import org.slf4j.Logger; import org.slf4j.LoggerFactory; - import org.springframework.context.ApplicationEventPublisher; import org.springframework.context.ApplicationEventPublisherAware; import org.springframework.dao.OptimisticLockingFailureException; @@ -529,14 +527,22 @@ public class CouchbaseTemplate implements CouchbaseOperations, ApplicationEventP public Boolean doInBucket() throws InterruptedException, ExecutionException { Document doc = encodeAndWrap(converted, version); Document storedDoc; - boolean checkVersion = version != null && version > 0L; + //We will check version only if required + boolean versionPresent = versionProperty != null; + //If version is not set - assumption that document is new, otherwise updating + boolean existingDocument = version != null && version > 0L; try { switch (persistType) { case SAVE: - if (checkVersion) { + if (!versionPresent) { + //No version field - no cas + storedDoc = client.upsert(doc, persistTo, replicateTo); + } else if (existingDocument) { + //Updating existing document with cas storedDoc = client.replace(doc, persistTo, replicateTo); } else { - storedDoc = client.upsert(doc, persistTo, replicateTo); + //Creating new document + storedDoc = client.insert(doc, persistTo, replicateTo); } break; case UPDATE: @@ -554,6 +560,9 @@ public class CouchbaseTemplate implements CouchbaseOperations, ApplicationEventP return true; } return false; + } catch (DocumentAlreadyExistsException e) { + throw new OptimisticLockingFailureException(persistType.getSpringDataOperationName() + + " document with version value failed: " + version, e); } catch (CASMismatchException e) { throw new OptimisticLockingFailureException(persistType.getSpringDataOperationName() + " document with version value failed: " + version, e);