From c9ba1e23653dc56cca34afc0f851bc71ba168783 Mon Sep 17 00:00:00 2001 From: "J. Brisbin" Date: Thu, 2 Dec 2010 09:56:32 -0600 Subject: [PATCH] Added QosParameters --- .../riak/core/KeyValueStoreOperations.java | 29 +++++++++++ .../data/riak/core/QosParameters.java | 34 +++++++++++++ .../data/riak/core/RiakQosParameters.java | 39 ++++++++++++++ .../data/riak/core/RiakTemplate.java | 51 +++++++++++++++++-- .../data/riak/core/RiakTemplateSpec.groovy | 15 ++++++ 5 files changed, 164 insertions(+), 4 deletions(-) create mode 100644 spring-data-riak/src/main/java/org/springframework/data/riak/core/QosParameters.java create mode 100644 spring-data-riak/src/main/java/org/springframework/data/riak/core/RiakQosParameters.java diff --git a/spring-data-riak/src/main/java/org/springframework/data/riak/core/KeyValueStoreOperations.java b/spring-data-riak/src/main/java/org/springframework/data/riak/core/KeyValueStoreOperations.java index b746b700e..799c6540f 100644 --- a/spring-data-riak/src/main/java/org/springframework/data/riak/core/KeyValueStoreOperations.java +++ b/spring-data-riak/src/main/java/org/springframework/data/riak/core/KeyValueStoreOperations.java @@ -35,6 +35,18 @@ public interface KeyValueStoreOperations { */ KeyValueStoreOperations set(K key, V value); + /** + * Variation on set() that allows the user to specify {@link org.springframework.data.riak.core.QosParameters}. + * + * @param key + * @param value + * @param qosParams + * @param + * @param + * @return + */ + KeyValueStoreOperations set(K key, V value, QosParameters qosParams); + /** * Set a value as a byte array at a specified key. * @@ -44,6 +56,20 @@ public interface KeyValueStoreOperations { */ KeyValueStoreOperations setAsBytes(K key, byte[] value); + /** + * Variation on setWithMetaData() that allows the user to pass {@link + * org.springframework.data.riak.core.QosParameters}. + * + * @param key + * @param value + * @param metaData + * @param qosParams + * @param + * @param + * @return + */ + KeyValueStoreOperations setWithMetaData(K key, V value, Map metaData, QosParameters qosParams); + // Get operations /** @@ -237,4 +263,7 @@ public interface KeyValueStoreOperations { */ Map getBucketSchema(B bucket, boolean listKeys); + KeyValueStoreOperations setAsBytes(K key, byte[] value, QosParameters qosParams); + + KeyValueStoreOperations setWithMetaData(K key, V value, Map metaData); } diff --git a/spring-data-riak/src/main/java/org/springframework/data/riak/core/QosParameters.java b/spring-data-riak/src/main/java/org/springframework/data/riak/core/QosParameters.java new file mode 100644 index 000000000..f69c4ec02 --- /dev/null +++ b/spring-data-riak/src/main/java/org/springframework/data/riak/core/QosParameters.java @@ -0,0 +1,34 @@ +package org.springframework.data.riak.core; + +/** + * Specify Quality Of Service parameters. + * + * @author J. Brisbin + */ +public interface QosParameters { + + /** + * Instruct the server on the read threshold. + * + * @param + * @return + */ + public T getReadThreshold(); + + /** + * Instruct the server on the normal write threshold. + * + * @param + * @return + */ + public T getWriteThreshold(); + + /** + * Instruct the server on the durable write threshold. + * + * @param + * @return + */ + public T getDurableWriteThreshold(); + +} diff --git a/spring-data-riak/src/main/java/org/springframework/data/riak/core/RiakQosParameters.java b/spring-data-riak/src/main/java/org/springframework/data/riak/core/RiakQosParameters.java new file mode 100644 index 000000000..8fbf14736 --- /dev/null +++ b/spring-data-riak/src/main/java/org/springframework/data/riak/core/RiakQosParameters.java @@ -0,0 +1,39 @@ +package org.springframework.data.riak.core; + +/** + * A generic class for specifying Quality Of Service parameters on operations. + * + * @author J. Brisbin + */ +@SuppressWarnings({"unchecked"}) +public class RiakQosParameters implements QosParameters { + + public Object readThreshold = null; + public Object writeThreshold = null; + public Object durableWriteThreshold = null; + + public void setReadThreshold(T readThreshold) { + this.readThreshold = readThreshold; + } + + public void setWriteThreshold(T writeThreshold) { + this.writeThreshold = writeThreshold; + } + + public void setDurableWriteThreshold(T durableWriteThreshold) { + this.durableWriteThreshold = durableWriteThreshold; + } + + public T getReadThreshold() { + return (T) this.readThreshold; + } + + public T getWriteThreshold() { + return (T) this.writeThreshold; + } + + public T getDurableWriteThreshold() { + return (T) this.durableWriteThreshold; + } + +} diff --git a/spring-data-riak/src/main/java/org/springframework/data/riak/core/RiakTemplate.java b/spring-data-riak/src/main/java/org/springframework/data/riak/core/RiakTemplate.java index 665af5626..0a7b2a578 100644 --- a/spring-data-riak/src/main/java/org/springframework/data/riak/core/RiakTemplate.java +++ b/spring-data-riak/src/main/java/org/springframework/data/riak/core/RiakTemplate.java @@ -141,6 +141,10 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe * A list of resolvers to turn a single object into a {@link BucketKeyPair}. */ protected List bucketKeyResolvers; + /** + * The default QosParameters to use for all operations through this template. + */ + protected QosParameters defaultQosParameters = null; /** * Take all the defaults. @@ -156,6 +160,7 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe */ public RiakTemplate(ClientHttpRequestFactory requestFactory) { super(requestFactory); + setRestTemplate(new RestTemplate()); } /** @@ -234,17 +239,27 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe return setWithMetaData(key, value, null); } + public KeyValueStoreOperations set(K key, V value, QosParameters qosParams) { + return setWithMetaData(key, value, null, qosParams); + } + public KeyValueStoreOperations setAsBytes(K key, byte[] value) { + return setAsBytes(key, value, null); + } + + public KeyValueStoreOperations setAsBytes(K key, byte[] value, QosParameters qosParams) { Assert.notNull(key, "Can't store an object with a NULL key."); BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, value); String bucketName = (null != bucketKeyPair.getBucket() ? bucketKeyPair.getBucket() .toString() : "bytes"); + String keyName = (null != qosParams ? bucketKeyPair.getKey() + .toString() + extractQosParameters(qosParams) : bucketKeyPair.getKey().toString()); RestTemplate restTemplate = getRestTemplate(); HttpHeaders headers = new HttpHeaders(); headers.set("X-Riak-ClientId", RIAK_CLIENT_ID); headers.setContentType(MediaType.APPLICATION_OCTET_STREAM); HttpEntity entity = new HttpEntity(value, headers); - restTemplate.put(defaultUri, entity, bucketName, bucketKeyPair.getKey()); + restTemplate.put(defaultUri, entity, bucketName, keyName); if (log.isDebugEnabled()) { log.debug(String.format("PUT byte[]: bucket=%s, key=%s", bucketKeyPair.getBucket(), @@ -253,8 +268,10 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe return this; } - public KeyValueStoreOperations setWithMetaData(K key, V value, Map metaData) { + public KeyValueStoreOperations setWithMetaData(K key, V value, Map metaData, QosParameters qosParams) { BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, value); + String keyName = (null != qosParams ? bucketKeyPair.getKey() + .toString() + extractQosParameters(qosParams) : bucketKeyPair.getKey().toString()); RestTemplate restTemplate = getRestTemplate(); HttpHeaders headers = new HttpHeaders(); headers.set("X-Riak-ClientId", RIAK_CLIENT_ID); @@ -268,7 +285,7 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe restTemplate.put(defaultUri, entity, bucketKeyPair.getBucket(), - bucketKeyPair.getKey()); + keyName); if (log.isDebugEnabled()) { log.debug(String.format("PUT object: bucket=%s, key=%s, value=%s", bucketKeyPair.getBucket(), @@ -278,6 +295,10 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe return this; } + public KeyValueStoreOperations setWithMetaData(K key, V value, Map metaData) { + return setWithMetaData(key, value, metaData, null); + } + /*----------------- Get Operations -----------------*/ public RiakValue getWithMetaData(K key, Class requiredType) { @@ -872,7 +893,6 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe return meta; } - protected T checkCache(K key, Class requiredType) { BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, requiredType); RiakValue obj = cache.get(bucketKeyPair); @@ -903,4 +923,27 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe } } + protected String extractQosParameters(QosParameters qosParams) { + List params = new LinkedList(); + if (null != qosParams.getReadThreshold()) { + params.add(String.format("r=%s", qosParams.getReadThreshold())); + } else if (null != defaultQosParameters && null != defaultQosParameters.getReadThreshold()) { + params.add(String.format("r=%s", defaultQosParameters.getReadThreshold())); + } + if (null != qosParams.getWriteThreshold()) { + params.add(String.format("w=%s", qosParams.getWriteThreshold())); + } else if (null != defaultQosParameters && null != defaultQosParameters.getWriteThreshold()) { + params.add(String.format("w=%s", defaultQosParameters.getWriteThreshold())); + } + if (null != qosParams.getDurableWriteThreshold()) { + params.add(String.format("dw=%s", qosParams.getDurableWriteThreshold())); + } else if (null != defaultQosParameters && null != defaultQosParameters.getDurableWriteThreshold()) { + params.add(String.format("dw=%s", defaultQosParameters.getDurableWriteThreshold())); + } + + return (params.size() > 0 ? "?" + StringUtils.collectionToDelimitedString( + params, + "&") : ""); + } + } diff --git a/spring-data-riak/src/test/groovy/org/springframework/data/riak/core/RiakTemplateSpec.groovy b/spring-data-riak/src/test/groovy/org/springframework/data/riak/core/RiakTemplateSpec.groovy index 3371eeffd..8b7ac02ad 100644 --- a/spring-data-riak/src/test/groovy/org/springframework/data/riak/core/RiakTemplateSpec.groovy +++ b/spring-data-riak/src/test/groovy/org/springframework/data/riak/core/RiakTemplateSpec.groovy @@ -94,6 +94,21 @@ class RiakTemplateSpec extends Specification { } + def "Test setting QosParameters"() { + + given: + def obj = riak.get("test:test") + + when: + def qos = new RiakQosParameters() + qos.durableWriteThreshold = "all" + riak.set("test:test", obj, qos) + + then: + true + + } + def "Test containsKey"() { when: