Bug fixes, added link(), started on linkWalk()
This commit is contained in:
@@ -45,6 +45,7 @@ import org.springframework.web.client.support.RestGatewaySupport;
|
||||
import java.io.ByteArrayOutputStream;
|
||||
import java.io.IOException;
|
||||
import java.io.InputStream;
|
||||
import java.io.StringWriter;
|
||||
import java.lang.annotation.Annotation;
|
||||
import java.text.ParseException;
|
||||
import java.text.SimpleDateFormat;
|
||||
@@ -144,20 +145,7 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
|
||||
/*----------------- Set Operations -----------------*/
|
||||
|
||||
public <K, V> KeyValueStoreOperations set(K key, V value) {
|
||||
BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, value);
|
||||
RestTemplate restTemplate = getRestTemplate();
|
||||
HttpHeaders headers = new HttpHeaders();
|
||||
headers.set("X-Riak-ClientId", RIAK_CLIENT_ID);
|
||||
headers.setContentType(extractMediaType(value));
|
||||
HttpEntity<V> entity = new HttpEntity<V>(value, headers);
|
||||
restTemplate.put(defaultUri, entity, bucketKeyPair.getBucket(), bucketKeyPair.getKey());
|
||||
if (log.isDebugEnabled()) {
|
||||
log.debug(String.format("PUT object: bucket=%s, key=%s, value=%s",
|
||||
bucketKeyPair.getBucket(),
|
||||
bucketKeyPair.getKey(),
|
||||
value));
|
||||
}
|
||||
return this;
|
||||
return setWithMetaData(key, value, null);
|
||||
}
|
||||
|
||||
public <K> KeyValueStoreOperations setAsBytes(K key, byte[] value) {
|
||||
@@ -176,6 +164,28 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
|
||||
return this;
|
||||
}
|
||||
|
||||
public <K, V> KeyValueStoreOperations setWithMetaData(K key, V value, Map<String, String> metaData) {
|
||||
BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, value);
|
||||
RestTemplate restTemplate = getRestTemplate();
|
||||
HttpHeaders headers = new HttpHeaders();
|
||||
headers.set("X-Riak-ClientId", RIAK_CLIENT_ID);
|
||||
headers.setContentType(extractMediaType(value));
|
||||
if (null != metaData) {
|
||||
for (Map.Entry<String, String> entry : metaData.entrySet()) {
|
||||
headers.set(entry.getKey(), entry.getValue());
|
||||
}
|
||||
}
|
||||
HttpEntity<V> entity = new HttpEntity<V>(value, headers);
|
||||
restTemplate.put(defaultUri, entity, bucketKeyPair.getBucket(), bucketKeyPair.getKey());
|
||||
if (log.isDebugEnabled()) {
|
||||
log.debug(String.format("PUT object: bucket=%s, key=%s, value=%s",
|
||||
bucketKeyPair.getBucket(),
|
||||
bucketKeyPair.getKey(),
|
||||
value));
|
||||
}
|
||||
return this;
|
||||
}
|
||||
|
||||
/*----------------- Get Operations -----------------*/
|
||||
|
||||
public <K, T> RiakValue<T> getWithMetaData(K key, Class<T> requiredType) {
|
||||
@@ -221,11 +231,13 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
|
||||
} catch (Throwable ignored) {
|
||||
targetClass = Map.class;
|
||||
}
|
||||
return (V) getWithMetaData(bucketKeyPair, targetClass).get();
|
||||
RiakValue<V> obj = getWithMetaData(bucketKeyPair, targetClass);
|
||||
return (null != obj ? obj.get() : null);
|
||||
}
|
||||
|
||||
public <K> byte[] getAsBytes(K key) {
|
||||
return getAsBytesWithMetaData(key).get();
|
||||
RiakValue<byte[]> obj = getAsBytesWithMetaData(key);
|
||||
return (null != obj ? obj.get() : null);
|
||||
}
|
||||
|
||||
public <K> RiakValue<byte[]> getAsBytesWithMetaData(K key) {
|
||||
@@ -283,7 +295,8 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
|
||||
return (T) obj;
|
||||
}
|
||||
}
|
||||
return getWithMetaData(key, requiredType).get();
|
||||
RiakValue<T> obj = getWithMetaData(key, requiredType);
|
||||
return (null != obj ? obj.get() : null);
|
||||
}
|
||||
|
||||
public <K, V> V getAndSet(K key, V value) {
|
||||
@@ -410,9 +423,9 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
|
||||
throw new DataAccessResourceFailureException(e.getMessage(), e);
|
||||
}
|
||||
}
|
||||
if (!stillExists) {
|
||||
stillExists = containsKey(key);
|
||||
}
|
||||
//if (!stillExists) {
|
||||
//stillExists = containsKey(key);
|
||||
//}
|
||||
}
|
||||
return !stillExists;
|
||||
}
|
||||
@@ -448,6 +461,9 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
|
||||
RestTemplate restTemplate = getRestTemplate();
|
||||
|
||||
RiakValue<byte[]> fromObj = getAsBytesWithMetaData(source);
|
||||
if (null == fromObj) {
|
||||
throw new DataStoreOperationException("Cannot link from a non-existent source: " + source);
|
||||
}
|
||||
HttpHeaders headers = new HttpHeaders();
|
||||
headers.setContentType(fromObj.getMetaData().getContentType());
|
||||
Object linksObj = fromObj.getMetaData().getProperties().get("Link");
|
||||
@@ -458,15 +474,45 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
|
||||
links.add(linksObj.toString());
|
||||
}
|
||||
links.add(String.format("<%s/%s/%s>; riaktag=\"%s\"", extractPrefix(), bkpTo.getBucket(), bkpTo.getKey(), tag));
|
||||
StringWriter sw = new StringWriter();
|
||||
boolean needsComma = false;
|
||||
for (String link : links) {
|
||||
headers.set("Link", link);
|
||||
if (!sw.toString().contains(link)) {
|
||||
if (needsComma) {
|
||||
sw.write(", ");
|
||||
} else {
|
||||
needsComma = true;
|
||||
}
|
||||
sw.write(link);
|
||||
}
|
||||
}
|
||||
headers.set("Link", sw.toString());
|
||||
HttpEntity entity = new HttpEntity(fromObj.get(), headers);
|
||||
restTemplate.put(defaultUri, entity, bkpFrom.getBucket(), bkpFrom.getKey());
|
||||
|
||||
return this;
|
||||
}
|
||||
|
||||
public <T, K> T linkWalk(K source, String tag) {
|
||||
BucketKeyPair bkpSource = resolveBucketKeyPair(source, null);
|
||||
RestTemplate restTemplate = getRestTemplate();
|
||||
final List<MediaType> types = new ArrayList<MediaType>();
|
||||
types.add(MediaType.ALL);
|
||||
restTemplate.execute(defaultUri + "/_,{tag},_", HttpMethod.GET, new RequestCallback() {
|
||||
public void doWithRequest(ClientHttpRequest request) throws IOException {
|
||||
request.getHeaders().setAccept(types);
|
||||
}
|
||||
}, new ResponseExtractor<Object>() {
|
||||
public Object extractData(ClientHttpResponse response) throws IOException {
|
||||
response.getHeaders();
|
||||
return null; //To change body of implemented methods use File | Settings | File Templates.
|
||||
}
|
||||
}, bkpSource.getBucket(),
|
||||
bkpSource.getKey(),
|
||||
tag);
|
||||
return null;
|
||||
}
|
||||
|
||||
/*----------------- Bucket Operations -----------------*/
|
||||
|
||||
public <B> Map<String, Object> getBucketSchema(B bucket) {
|
||||
@@ -493,13 +539,13 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
|
||||
bucketKeyResolvers.add(new SimpleBucketKeyResolver());
|
||||
}
|
||||
|
||||
List<HttpMessageConverter<?>> converters = getRestTemplate().getMessageConverters();
|
||||
if (groovyPresent) {
|
||||
// Native conversion for Groovy GString objects
|
||||
ObjectMapper mapper = new ObjectMapper();
|
||||
CustomSerializerFactory fac = new CustomSerializerFactory();
|
||||
fac.addSpecificMapping(GStringImpl.class, ToStringSerializer.instance);
|
||||
mapper.setSerializerFactory(fac);
|
||||
List<HttpMessageConverter<?>> converters = getRestTemplate().getMessageConverters();
|
||||
for (HttpMessageConverter converter : converters) {
|
||||
if (converter instanceof MappingJacksonHttpMessageConverter) {
|
||||
((MappingJacksonHttpMessageConverter) converter).setObjectMapper(mapper);
|
||||
|
||||
@@ -31,4 +31,9 @@ public class SimpleBucketKeyPair<B, K> implements BucketKeyPair, Comparable {
|
||||
}
|
||||
return -1;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return String.format("{bucket=%s, key=%s}", bucket, key);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user