+ * Portions (c) 2010 by NPC International, Inc. or the
+ * original author(s).
+ *
+ * Licensed under the Apache License, Version 2.0 (the "License");
+ * you may not use this file except in compliance with the License.
+ * You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.springframework.data.keyvalue.riak.core;
+
+import org.springframework.beans.factory.InitializingBean;
+import org.springframework.data.keyvalue.riak.mapreduce.MapReduceJob;
+import org.springframework.data.keyvalue.riak.mapreduce.MapReduceOperations;
+import org.springframework.data.keyvalue.riak.mapreduce.RiakMapReduceJob;
+import org.springframework.http.client.ClientHttpRequestFactory;
+import org.springframework.util.Assert;
+import org.springframework.web.client.RestTemplate;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.Future;
+
+/**
+ * An implementation of {@link org.springframework.data.keyvalue.riak.core.KeyValueStoreOperations}
+ * and {@link org.springframework.data.keyvalue.riak.mapreduce.MapReduceOperations} for the Riak
+ * data store.
+ *
+ * To use the RiakTemplate, create a singleton in your Spring application-context.xml:
+ *
+ * <bean id="riak" class="org.springframework.data.keyvalue.riak.core.RiakTemplate"
+ * p:defaultUri="http://localhost:8098/riak/{bucket}/{key}"
+ * p:mapReduceUri="http://localhost:8098/mapred"/>
+ *
+ * To store and retrieve objects in Riak, use the setXXX and getXXX methods (example in
+ * Groovy):
+ *
+ * def obj = new TestObject(name: "My Name", age: 40)
+ * riak.set([bucket: "mybucket", key: "mykey"], obj)
+ * ...
+ * def name = riak.get([bucket: "mybucket", key: "mykey"]).name
+ * println "Hello $name!"
+ *
+ * You're key object should be one of: - A
String encoding the bucket and key
+ * together, separated by a colon. e.g. "mybucket:mykey" - An implementation of
+ * BucketKeyPair (like {@link org.springframework.data.keyvalue.riak.core.SimpleBucketKeyPair})
+ * - A
Map with both a "bucket" and a "key" specified. - A
+ *
String of only the key name, but specifying a bucket by using the {@link
+ * org.springframework.data.keyvalue.riak.convert.KeyValueStoreMetaData} annotation on the
+ * object you're storing.
+ *
+ * @author J. Brisbin
+ */
+public class RiakKeyValueTemplate extends AbstractRiakTemplate implements KeyValueStoreOperations, MapReduceOperations, InitializingBean {
+
+ protected RiakTemplate riak;
+
+ /**
+ * Take all the defaults.
+ */
+ public RiakKeyValueTemplate() {
+ super();
+ riak = new RiakTemplate();
+ }
+
+ /**
+ * Use the specified {@link org.springframework.http.client.ClientHttpRequestFactory}.
+ *
+ * @param requestFactory
+ */
+ public RiakKeyValueTemplate(ClientHttpRequestFactory requestFactory) {
+ super(requestFactory);
+ riak = new RiakTemplate(requestFactory);
+ }
+
+ /**
+ * Use the specified defaultUri and mapReduceUri.
+ *
+ * @param defaultUri
+ * @param mapReduceUri
+ */
+ public RiakKeyValueTemplate(String defaultUri, String mapReduceUri) {
+ setRestTemplate(new RestTemplate());
+ this.setDefaultUri(defaultUri);
+ this.mapReduceUri = mapReduceUri;
+ this.riak = new RiakTemplate(defaultUri, mapReduceUri);
+ }
+
+ @Override
+ public void afterPropertiesSet() throws Exception {
+ super.afterPropertiesSet();
+ riak.afterPropertiesSet();
+ }
+
+ /*----------------- Set Operations -----------------*/
+
+ public KeyValueStoreOperations set(K key, V value) {
+ return setWithMetaData(key, value, null, 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, "Key cannot be null!");
+ BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, value);
+ riak.setAsBytes(bucketKeyPair.getBucket(), bucketKeyPair.getKey(), value, qosParams);
+ return this;
+ }
+
+ public KeyValueStoreOperations setWithMetaData(K key, V value, Map metaData, QosParameters qosParams) {
+ BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, value);
+ riak.setWithMetaData(bucketKeyPair.getBucket(),
+ bucketKeyPair.getKey(),
+ value,
+ metaData,
+ qosParams);
+ return this;
+ }
+
+ public KeyValueStoreOperations setWithMetaData(K key, V value, Map metaData) {
+ BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, value);
+ riak.setWithMetaData(bucketKeyPair.getBucket(),
+ bucketKeyPair.getKey(),
+ value,
+ metaData,
+ null);
+ return this;
+ }
+
+ /*----------------- Get Operations -----------------*/
+
+ public RiakMetaData getMetaData(K key) {
+ BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, null);
+ return riak.getMetaData(bucketKeyPair.getBucket(), bucketKeyPair.getKey());
+ }
+
+ public RiakValue getWithMetaData(K key, Class requiredType) {
+ BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, null);
+ return riak.getWithMetaData(bucketKeyPair.getBucket(),
+ bucketKeyPair.getKey(),
+ requiredType);
+ }
+
+ @SuppressWarnings({"unchecked"})
+ public V get(K key) {
+ BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, null);
+ return (V) riak.get(bucketKeyPair.getBucket(), bucketKeyPair.getKey());
+ }
+
+ public byte[] getAsBytes(K key) {
+ BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, null);
+ RiakValue obj = riak.getAsBytesWithMetaData(bucketKeyPair.getBucket(),
+ bucketKeyPair.getKey());
+ return (null != obj ? obj.get() : null);
+ }
+
+ public RiakValue getAsBytesWithMetaData(K key) {
+ BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, null);
+ return riak.getAsBytesWithMetaData(bucketKeyPair.getBucket(), bucketKeyPair.getKey());
+ }
+
+ public T getAsType(K key, Class requiredType) {
+ BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, null);
+ return riak.getAsType(bucketKeyPair.getBucket(), bucketKeyPair.getKey(), requiredType);
+ }
+
+ public V getAndSet(K key, V value) {
+ BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, null);
+ return riak.getAndSet(bucketKeyPair.getBucket(), bucketKeyPair.getKey(), value);
+ }
+
+ public byte[] getAndSetAsBytes(K key, byte[] value) {
+ BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, null);
+ return riak.getAndSetAsBytes(bucketKeyPair.getBucket(), bucketKeyPair.getKey(), value);
+ }
+
+ public T getAndSetAsType(K key, V value, Class requiredType) {
+ BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, null);
+ return riak.getAndSetAsType(bucketKeyPair.getBucket(),
+ bucketKeyPair.getKey(),
+ value,
+ requiredType);
+ }
+
+ @SuppressWarnings({"unchecked"})
+ public List getValues(List keys) {
+ List results = new ArrayList();
+ for (K key : keys) {
+ BucketKeyPair bkp = resolveBucketKeyPair(key, null);
+ results.add((V) riak.get(bkp.getBucket(), bkp.getKey()));
+ }
+ return results;
+ }
+
+ public List getValues(K... keys) {
+ return getValues(keys);
+ }
+
+ public List getValuesAsType(List keys, Class requiredType) {
+ List results = new ArrayList();
+ for (K key : keys) {
+ BucketKeyPair bkp = resolveBucketKeyPair(key, null);
+ results.add(riak.getAsType(bkp.getBucket(), bkp.getKey(), requiredType));
+ }
+ return results;
+ }
+
+ public List getValuesAsType(Class requiredType, K... keys) {
+ return riak.getValuesAsType(requiredType, keys);
+ }
+
+ /*----------------- Only-Set-Once Operations -----------------*/
+
+ public KeyValueStoreOperations setIfKeyNonExistent(K key, V value) {
+ BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, null);
+ riak.setIfKeyNonExistent(bucketKeyPair.getBucket(), bucketKeyPair.getKey(), value);
+ return this;
+ }
+
+ public KeyValueStoreOperations setIfKeyNonExistentAsBytes(K key, byte[] value) {
+ BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, null);
+ riak.setIfKeyNonExistent(bucketKeyPair.getBucket(), bucketKeyPair.getKey(), value);
+ return this;
+ }
+
+ /*----------------- Multiple Item Operations -----------------*/
+
+ public KeyValueStoreOperations setMultiple(Map keysAndValues) {
+ for (Map.Entry entry : keysAndValues.entrySet()) {
+ set(entry.getKey(), entry.getValue());
+ }
+ return this;
+ }
+
+ public KeyValueStoreOperations setMultipleAsBytes(Map keysAndValues) {
+ for (Map.Entry entry : keysAndValues.entrySet()) {
+ setAsBytes(entry.getKey(), entry.getValue());
+ }
+ return this;
+ }
+
+ public KeyValueStoreOperations setMultipleIfKeysNonExistent(Map keysAndValues) {
+ for (Map.Entry entry : keysAndValues.entrySet()) {
+ setIfKeyNonExistent(entry.getKey(), entry.getValue());
+ }
+ return this;
+ }
+
+ public KeyValueStoreOperations setMultipleAsBytesIfKeysNonExistent(Map keysAndValues) {
+ for (Map.Entry entry : keysAndValues.entrySet()) {
+ setIfKeyNonExistentAsBytes(entry.getKey(), entry.getValue());
+ }
+ return this;
+ }
+
+ /*----------------- Key Operations -----------------*/
+
+ public boolean containsKey(K key) {
+ BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, null);
+ return riak.containsKey(bucketKeyPair.getBucket(), bucketKeyPair.getKey());
+ }
+
+ public boolean deleteKeys(K... keys) {
+ return riak.deleteKeys(keys);
+ }
+
+ /*----------------- Map/Reduce Operations -----------------*/
+
+ public RiakMapReduceJob createMapReduceJob() {
+ return new RiakMapReduceJob(riak);
+ }
+
+ public Object execute(MapReduceJob job) {
+ return execute(job, List.class);
+ }
+
+ public T execute(MapReduceJob job, Class targetType) {
+ return riak.execute(job, targetType);
+ }
+
+ public Future> submit(MapReduceJob job) {
+ // Run this job asynchronously.
+ return riak.submit(job);
+ }
+
+ /*----------------- Link Operations -----------------*/
+
+ /**
+ * Use Riak's native Link mechanism to link two entries together.
+ *
+ * @param destination Key to the child object
+ * @param source Key to the parent object
+ * @param tag The tag for this relationship
+ * @return This template interface
+ */
+ public RiakKeyValueTemplate link(K1 destination, K2 source, String tag) {
+ BucketKeyPair bkpFrom = resolveBucketKeyPair(source, null);
+ BucketKeyPair bkpTo = resolveBucketKeyPair(destination, null);
+ riak.link(bkpTo.getBucket(), bkpTo.getKey(), bkpFrom.getBucket(), bkpFrom.getKey(), tag);
+ return this;
+ }
+
+ /**
+ * Use Riak's link walking mechanism to retrieve a multipart message that will be decoded like
+ * they were individual objects (e.g. using the built-in HttpMessageConverters of
+ * RestTemplate).
+ *
+ * @param source
+ * @param tag
+ * @return
+ */
+ @SuppressWarnings({"unchecked"})
+ public T linkWalk(K source, String tag) {
+ BucketKeyPair bkpSource = resolveBucketKeyPair(source, null);
+ return (T) riak.linkWalk(bkpSource.getBucket(), bkpSource.getKey(), tag);
+ }
+
+ /*----------------- Bucket Operations -----------------*/
+
+ public Map getBucketSchema(B bucket) {
+ return riak.getBucketSchema(bucket, false);
+ }
+
+ public Map getBucketSchema(B bucket, boolean listKeys) {
+ return riak.getBucketSchema(bucket, listKeys);
+ }
+
+ public KeyValueStoreOperations updateBucketSchema(B bucket, Map props) {
+ riak.updateBucketSchema(bucket, props);
+ return this;
+ }
+
+}
diff --git a/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/RiakMetaData.java b/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/RiakMetaData.java
index 69d3a6cd2..2c96eff21 100644
--- a/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/RiakMetaData.java
+++ b/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/RiakMetaData.java
@@ -20,6 +20,7 @@ package org.springframework.data.keyvalue.riak.core;
import org.springframework.http.MediaType;
+import java.util.Date;
import java.util.Map;
/**
@@ -46,6 +47,10 @@ public class RiakMetaData implements KeyValueStoreMetaData {
return mediaType;
}
+ public long getLastModified() {
+ return ((Date) properties.get("Last-Modified")).getTime();
+ }
+
public Map getProperties() {
return this.properties;
}
diff --git a/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/RiakTemplate.java b/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/RiakTemplate.java
index 723acd6d1..ec7d1df40 100644
--- a/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/RiakTemplate.java
+++ b/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/RiakTemplate.java
@@ -18,18 +18,9 @@
package org.springframework.data.keyvalue.riak.core;
-import org.codehaus.groovy.runtime.GStringImpl;
-import org.codehaus.jackson.map.ObjectMapper;
-import org.codehaus.jackson.map.ser.CustomSerializerFactory;
-import org.codehaus.jackson.map.ser.ToStringSerializer;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-import org.springframework.beans.factory.InitializingBean;
import org.springframework.core.convert.ConversionService;
-import org.springframework.core.convert.support.ConversionServiceFactory;
import org.springframework.dao.DataAccessResourceFailureException;
import org.springframework.data.keyvalue.riak.DataStoreOperationException;
-import org.springframework.data.keyvalue.riak.convert.KeyValueStoreMetaData;
import org.springframework.data.keyvalue.riak.mapreduce.MapReduceJob;
import org.springframework.data.keyvalue.riak.mapreduce.MapReduceOperations;
import org.springframework.data.keyvalue.riak.mapreduce.RiakMapReduceJob;
@@ -38,12 +29,9 @@ import org.springframework.http.client.ClientHttpRequest;
import org.springframework.http.client.ClientHttpRequestFactory;
import org.springframework.http.client.ClientHttpResponse;
import org.springframework.http.converter.HttpMessageConverter;
-import org.springframework.http.converter.json.MappingJacksonHttpMessageConverter;
import org.springframework.util.Assert;
-import org.springframework.util.ClassUtils;
import org.springframework.util.StringUtils;
import org.springframework.web.client.*;
-import org.springframework.web.client.support.RestGatewaySupport;
import javax.mail.BodyPart;
import javax.mail.MessagingException;
@@ -53,19 +41,11 @@ import java.io.ByteArrayOutputStream;
import java.io.EOFException;
import java.io.IOException;
import java.io.InputStream;
-import java.lang.annotation.Annotation;
-import java.text.ParseException;
-import java.text.SimpleDateFormat;
import java.util.*;
-import java.util.concurrent.ConcurrentSkipListMap;
-import java.util.concurrent.ExecutorService;
-import java.util.concurrent.Executors;
import java.util.concurrent.Future;
-import java.util.regex.Matcher;
-import java.util.regex.Pattern;
/**
- * An implementation of {@link org.springframework.data.keyvalue.riak.core.KeyValueStoreOperations}
+ * An implementation of {@link org.springframework.data.keyvalue.riak.core.BucketKeyValueStoreOperations}
* and {@link org.springframework.data.keyvalue.riak.mapreduce.MapReduceOperations} for the Riak
* data store.
*
@@ -79,84 +59,21 @@ import java.util.regex.Pattern;
* Groovy):
*
* def obj = new TestObject(name: "My Name", age: 40)
- * riak.set([bucket: "mybucket", key: "mykey"], obj)
+ * riak.set("mybucket", "mykey", obj)
* ...
- * def name = riak.get([bucket: "mybucket", key: "mykey"]).name
+ * def name = riak.get("mybucket", "mykey").name
* println "Hello $name!"
*
- * You're key object should be one of: - A
String encoding the bucket and key
- * together, separated by a colon. e.g. "mybucket:mykey" - An implementation of
- * BucketKeyPair (like {@link org.springframework.data.keyvalue.riak.core.SimpleBucketKeyPair})
- * - A
Map with both a "bucket" and a "key" specified. - A
- *
String of only the key name, but specifying a bucket by using the {@link
- * org.springframework.data.keyvalue.riak.convert.KeyValueStoreMetaData} annotation on the
- * object you're storing.
*
* @author J. Brisbin
*/
-@SuppressWarnings({"unchecked"})
-public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOperations, MapReduceOperations, InitializingBean {
-
- /**
- * Client ID used by Riak to correlate updates.
- */
- private static final String RIAK_CLIENT_ID = "org.springframework.data.keyvalue.riak.core.RiakTemplate/1.0";
- /**
- * Regex used to extract host, port, and prefix from the given URI.
- */
- private static final Pattern prefix = Pattern.compile(
- "http[s]?://(\\S+):([0-9]+)/(\\S+)/\\{bucket\\}(\\S+)");
- /**
- * Do we need to handle Groovy strings in the Jackson JSON processor?
- */
- private static final boolean groovyPresent = ClassUtils.isPresent(
- "org.codehaus.groovy.runtime.GStringImpl",
- RiakTemplate.class.getClassLoader());
- /**
- * For getting a java.util.Date from the Last-Modified header.
- */
- private static SimpleDateFormat httpDate = new SimpleDateFormat(
- "EEE, d MMM yyyy HH:mm:ss z");
-
- protected final Logger log = LoggerFactory.getLogger(getClass());
- /**
- * For converting objects to/from other kinds of objects.
- */
- protected ConversionService conversionService = ConversionServiceFactory.createDefaultConversionService();
- /**
- * For caching objects based on ETags.
- */
- protected ConcurrentSkipListMap> cache = new ConcurrentSkipListMap>();
- /**
- * Whether or not to use the ETag-based cache.
- */
- protected boolean useCache = true;
- /**
- * {@link ExecutorService} to use for running asynchronous jobs.
- */
- protected ExecutorService executorService = Executors.newCachedThreadPool();
- /**
- * The URI to use inside the RestTemplate.
- */
- protected String defaultUri = "http://localhost:8098/riak/{bucket}/{key}";
- /**
- * The URI for the Riak Map/Reduce API.
- */
- protected String mapReduceUri = "http://localhost:8098/mapred";
- /**
- * 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;
+public class RiakTemplate extends AbstractRiakTemplate implements BucketKeyValueStoreOperations, MapReduceOperations {
/**
* Take all the defaults.
*/
public RiakTemplate() {
- setRestTemplate(new RestTemplate());
+ super();
}
/**
@@ -166,7 +83,6 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
*/
public RiakTemplate(ClientHttpRequestFactory requestFactory) {
super(requestFactory);
- setRestTemplate(new RestTemplate());
}
/**
@@ -181,99 +97,27 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
this.mapReduceUri = mapReduceUri;
}
- public ConversionService getConversionService() {
- return conversionService;
- }
-
- /**
- * Specify the conversion service to use.
- *
- * @param conversionService
- */
- public void setConversionService(ConversionService conversionService) {
- this.conversionService = conversionService;
- }
-
- public String getDefaultUri() {
- return defaultUri;
- }
-
- public void setDefaultUri(String defaultUri) {
- this.defaultUri = defaultUri;
- }
-
- public String getMapReduceUri() {
- return mapReduceUri;
- }
-
- public void setMapReduceUri(String mapReduceUri) {
- this.mapReduceUri = mapReduceUri;
- }
-
- public List getBucketKeyResolvers() {
- return bucketKeyResolvers;
- }
-
- /**
- * Set the list of BucketKeyResolvers to use.
- *
- * @param bucketKeyResolvers
- */
- public void setBucketKeyResolvers(List bucketKeyResolvers) {
- this.bucketKeyResolvers = bucketKeyResolvers;
- }
-
- public boolean isUseCache() {
- return useCache;
- }
-
- public void setUseCache(boolean useCache) {
- this.useCache = useCache;
- }
-
- /**
- * Extract the prefix from the URI for use in creating links.
- *
- * @return
- */
- public String getPrefix() {
- Matcher m = prefix.matcher(defaultUri);
- if (m.matches()) {
- return "/" + m.group(3);
- }
- return "/riak";
- }
-
- public ExecutorService getExecutorService() {
- return executorService;
- }
-
- public void setExecutorService(ExecutorService executorService) {
- this.executorService = executorService;
- }
/*----------------- Set Operations -----------------*/
- public KeyValueStoreOperations set(K key, V value) {
- return setWithMetaData(key, value, null);
+ public BucketKeyValueStoreOperations set(B bucket, K key, V value) {
+ return setWithMetaData(bucket, key, value, null, null);
}
- public KeyValueStoreOperations set(K key, V value, QosParameters qosParams) {
- return setWithMetaData(key, value, null, qosParams);
+ public BucketKeyValueStoreOperations set(B bucket, K key, V value, QosParameters qosParams) {
+ return setWithMetaData(bucket, key, value, null, qosParams);
}
- public KeyValueStoreOperations setAsBytes(K key, byte[] value) {
- return setAsBytes(key, value, null);
+ public BucketKeyValueStoreOperations setAsBytes(B bucket, K key, byte[] value) {
+ return setAsBytes(bucket, 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);
+ public BucketKeyValueStoreOperations setAsBytes(B bucket, K key, byte[] value, QosParameters qosParams) {
+ Assert.notNull(key, "Key cannot be null!");
// If I don't give a bucket name, since I don't have an object type, use 'bytes'
- String bucketName = (null != bucketKeyPair.getBucket() ? bucketKeyPair.getBucket()
- .toString() : "bytes");
+ String bucketName = (null != bucket ? bucket.toString() : "bytes");
// Get a key name that may or may not include the QOS parameters.
- String keyName = (null != qosParams ? bucketKeyPair.getKey()
- .toString() + extractQosParameters(qosParams) : bucketKeyPair.getKey().toString());
+ String keyName = (null != qosParams ? key.toString() + extractQosParameters(qosParams) : key
+ .toString());
RestTemplate restTemplate = getRestTemplate();
HttpHeaders headers = new HttpHeaders();
headers.set("X-Riak-ClientId", RIAK_CLIENT_ID);
@@ -282,9 +126,7 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
try {
restTemplate.put(defaultUri, entity, bucketName, keyName);
if (log.isDebugEnabled()) {
- log.debug(String.format("PUT byte[]: bucket=%s, key=%s",
- bucketKeyPair.getBucket(),
- bucketKeyPair.getKey()));
+ log.debug(String.format("PUT byte[]: bucket=%s, key=%s", bucketName, keyName));
}
} catch (RestClientException e) {
throw new DataStoreOperationException(e.getMessage(), e);
@@ -292,11 +134,10 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
return this;
}
- public KeyValueStoreOperations setWithMetaData(K key, V value, Map metaData, QosParameters qosParams) {
- BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, value);
+ public BucketKeyValueStoreOperations setWithMetaData(B bucket, K key, V value, Map metaData, QosParameters qosParams) {
// Get a key name that may or may not include the QOS parameters.
- String keyName = (null != qosParams ? bucketKeyPair.getKey()
- .toString() + extractQosParameters(qosParams) : bucketKeyPair.getKey().toString());
+ String keyName = (null != qosParams ? key.toString() + extractQosParameters(qosParams) : key
+ .toString());
RestTemplate restTemplate = getRestTemplate();
HttpHeaders headers = new HttpHeaders();
headers.set("X-Riak-ClientId", RIAK_CLIENT_ID);
@@ -308,12 +149,9 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
}
HttpEntity entity = new HttpEntity(value, headers);
try {
- restTemplate.put(defaultUri, entity, bucketKeyPair.getBucket(), keyName);
+ restTemplate.put(defaultUri, entity, bucket, keyName);
if (log.isDebugEnabled()) {
- log.debug(String.format("PUT object: bucket=%s, key=%s, value=%s",
- bucketKeyPair.getBucket(),
- bucketKeyPair.getKey(),
- value));
+ log.debug(String.format("PUT object: bucket=%s, key=%s, value=%s", bucket, key, value));
}
} catch (RestClientException e) {
throw new DataStoreOperationException(e.getMessage(), e);
@@ -321,22 +159,33 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
return this;
}
- public KeyValueStoreOperations setWithMetaData(K key, V value, Map metaData) {
- return setWithMetaData(key, value, metaData, null);
+ public BucketKeyValueStoreOperations setWithMetaData(B bucket, K key, V value, Map metaData) {
+ return setWithMetaData(bucket, key, value, metaData, null);
}
/*----------------- Get Operations -----------------*/
- public RiakValue getWithMetaData(K key, Class requiredType) {
- BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, null);
+ public RiakMetaData getMetaData(B bucket, K key) {
+ RestTemplate restTemplate = getRestTemplate();
+ HttpHeaders headers = null;
+ try {
+ headers = restTemplate.headForHeaders(defaultUri, bucket, key);
+ return extractMetaData(headers);
+ } catch (ResourceAccessException e) {
+ } catch (IOException e) {
+ throw new DataAccessResourceFailureException(e.getMessage(), e);
+ }
+ return null;
+ }
+
+ public RiakValue getWithMetaData(B bucket, K key, Class requiredType) {
// If no bucket name is given, infer it from the type name.
- String bucketName = (null != bucketKeyPair.getBucket() ? bucketKeyPair.getBucket()
- .toString() : requiredType.getName());
+ String bucketName = (null != bucket ? bucket.toString() : requiredType.getName());
RestTemplate restTemplate = getRestTemplate();
if (log.isDebugEnabled()) {
log.debug(String.format("GET object: bucket=%s, key=%s, type=%s",
bucketName,
- bucketKeyPair.getKey(),
+ key,
requiredType.getName()));
}
@@ -344,12 +193,12 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
ResponseEntity result = restTemplate.getForEntity(defaultUri,
requiredType,
bucketName,
- bucketKeyPair.getKey());
+ key);
if (result.hasBody()) {
RiakMetaData meta = extractMetaData(result.getHeaders());
- RiakValue val = new RiakValue(result.getBody(), meta);
+ RiakValue val = new RiakValue(result.getBody(), meta);
if (useCache) {
- cache.put(bucketKeyPair, val);
+ cache.put(new SimpleBucketKeyPair