Support for linking objects, ETag-based caching

This commit is contained in:
J. Brisbin
2010-11-16 15:27:23 -06:00
parent 98b8756f24
commit 5436140834
10 changed files with 402 additions and 44 deletions

View File

@@ -0,0 +1,10 @@
package org.springframework.datastore.riak.core;
/**
* @author J. Brisbin <jon@jbrisbin.com>
*/
public interface BucketSchema {
String getName();
}

View File

@@ -0,0 +1,16 @@
package org.springframework.datastore.riak.core;
import org.springframework.http.MediaType;
import java.util.Map;
/**
* @author J. Brisbin <jon@jbrisbin.com>
*/
public interface KeyValueStoreMetaData {
MediaType getContentType();
Map<String, Object> getProperties();
}

View File

@@ -67,4 +67,8 @@ public interface KeyValueStoreOperations {
<K> boolean deleteKeys(K... keys);
<B> Map<String, Object> getBucketSchema(B bucket);
<B> Map<String, Object> getBucketSchema(B bucket, boolean listKeys);
}

View File

@@ -0,0 +1,12 @@
package org.springframework.datastore.riak.core;
/**
* @author J. Brisbin <jon@jbrisbin.com>
*/
public interface KeyValueStoreValue<T> {
KeyValueStoreMetaData getMetaData();
T get();
}

View File

@@ -0,0 +1,32 @@
package org.springframework.datastore.riak.core;
import org.springframework.http.MediaType;
import java.util.Map;
/**
* @author J. Brisbin <jon@jbrisbin.com>
*/
public class RiakMetaData implements KeyValueStoreMetaData {
private MediaType mediaType = MediaType.APPLICATION_JSON;
private Map<String, Object> properties;
public RiakMetaData(Map<String, Object> properties) {
this.properties = properties;
}
public RiakMetaData(MediaType mediaType, Map<String, Object> properties) {
this.mediaType = mediaType;
this.properties = properties;
}
public MediaType getContentType() {
return mediaType;
}
public Map<String, Object> getProperties() {
return this.properties;
}
}

View File

@@ -32,24 +32,32 @@ import org.springframework.datastore.riak.mapreduce.MapReduceJob;
import org.springframework.datastore.riak.mapreduce.MapReduceOperations;
import org.springframework.datastore.riak.mapreduce.RiakMapReduceJob;
import org.springframework.http.*;
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.web.client.HttpClientErrorException;
import org.springframework.web.client.ResourceAccessException;
import org.springframework.web.client.RestTemplate;
import org.springframework.web.client.*;
import org.springframework.web.client.support.RestGatewaySupport;
import java.io.ByteArrayOutputStream;
import java.io.IOException;
import java.io.InputStream;
import java.lang.annotation.Annotation;
import java.text.ParseException;
import java.text.SimpleDateFormat;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
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;
/**
* @author J. Brisbin <jon@jbrisbin.com>
@@ -57,12 +65,17 @@ import java.util.concurrent.Future;
@SuppressWarnings({"unchecked"})
public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOperations, MapReduceOperations, InitializingBean {
private static final String RIAK_CLIENT_ID = "org.springframework.datastore.riak.core.RiakTemplate/1.0";
private static final Pattern prefix = Pattern.compile("http[s]?://(\\S+):([0-9]+)/(\\S+)/\\{bucket\\}(\\S+)");
private static final boolean groovyPresent = ClassUtils.isPresent("org.codehaus.groovy.runtime.GStringImpl",
RiakTemplate.class.getClassLoader());
private static SimpleDateFormat httpDate = new SimpleDateFormat("EEE, d MMM yyyy HH:mm:ss z");
protected final Logger log = LoggerFactory.getLogger(getClass());
protected ConversionService conversionService = ConversionServiceFactory.createDefaultConversionService();
protected ConcurrentSkipListMap<Object, Object> cache = new ConcurrentSkipListMap<Object, Object>();
protected ObjectMapper mapper = new ObjectMapper();
protected ConcurrentSkipListMap<BucketKeyPair, RiakValue<?>> cache = new ConcurrentSkipListMap<BucketKeyPair, RiakValue<?>>();
protected boolean useCache = true;
protected ExecutorService queue = Executors.newCachedThreadPool();
protected String defaultUri = "http://localhost:8098/riak/{bucket}/{key}";
@@ -77,6 +90,17 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
super(requestFactory);
}
public RiakTemplate(String defaultUri) {
setRestTemplate(new RestTemplate());
setDefaultUri(defaultUri);
}
public RiakTemplate(String defaultUri, String mapReduceUri) {
setRestTemplate(new RestTemplate());
this.setDefaultUri(defaultUri);
this.mapReduceUri = mapReduceUri;
}
public ConversionService getConversionService() {
return conversionService;
}
@@ -109,10 +133,21 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
this.bucketKeyResolvers = bucketKeyResolvers;
}
public boolean isUseCache() {
return useCache;
}
public void setUseCache(boolean useCache) {
this.useCache = useCache;
}
/*----------------- 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());
@@ -131,6 +166,7 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
String bucketName = (null != bucketKeyPair.getBucket() ? bucketKeyPair.getBucket().toString() : "bytes");
RestTemplate restTemplate = getRestTemplate();
HttpHeaders headers = new HttpHeaders();
headers.set("X-Riak-ClientId", RIAK_CLIENT_ID);
headers.setContentType(MediaType.APPLICATION_OCTET_STREAM);
HttpEntity<byte[]> entity = new HttpEntity<byte[]>(value, headers);
restTemplate.put(defaultUri, entity, bucketName, bucketKeyPair.getKey());
@@ -140,36 +176,10 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
return this;
}
public <K, V> V get(K key) {
/*----------------- Get Operations -----------------*/
public <K, T> RiakValue<T> getWithMetaData(K key, Class<T> requiredType) {
BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, null);
RestTemplate restTemplate = getRestTemplate();
Class targetClass;
try {
targetClass = Class.forName(bucketKeyPair.getBucket().toString());
} catch (Throwable ignored) {
targetClass = Map.class;
}
String bucketName = (null != bucketKeyPair.getBucket() ? bucketKeyPair.getBucket()
.toString() : targetClass.getName());
if (log.isDebugEnabled()) {
log.debug(String.format("GET object: bucket=%s, key=%s", bucketName, bucketKeyPair.getKey()));
}
try {
return (V) restTemplate.getForObject(defaultUri, targetClass, bucketName, bucketKeyPair.getKey());
} catch (HttpClientErrorException e) {
if (e.getStatusCode() != HttpStatus.NOT_FOUND) {
throw new DataAccessResourceFailureException(e.getMessage(), e);
}
return null;
}
}
public <K> byte[] getAsBytes(K key) {
return getAsType(key, byte[].class);
}
public <K, T> T getAsType(K key, Class<T> requiredType) {
BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, requiredType);
String bucketName = (null != bucketKeyPair.getBucket() ? bucketKeyPair.getBucket()
.toString() : requiredType.getName());
RestTemplate restTemplate = getRestTemplate();
@@ -179,14 +189,101 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
bucketKeyPair.getKey(),
requiredType.getName()));
}
try {
return (T) restTemplate.getForObject(defaultUri, requiredType, bucketName, bucketKeyPair.getKey());
ResponseEntity<T> result = restTemplate.getForEntity(defaultUri,
requiredType,
bucketName,
bucketKeyPair.getKey());
if (result.hasBody()) {
RiakMetaData meta = extractMetaData(result.getHeaders());
RiakValue val = new RiakValue(result.getBody(), meta);
if (useCache) {
cache.put(bucketKeyPair, val);
}
return val;
}
} catch (HttpClientErrorException e) {
if (e.getStatusCode() != HttpStatus.NOT_FOUND) {
throw new DataAccessResourceFailureException(e.getMessage(), e);
throw new DataStoreOperationException(e.getMessage(), e);
}
return null;
} catch (IOException e) {
log.error(e.getMessage(), e);
}
return null;
}
public <K, V> V get(K key) {
BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, null);
Class targetClass;
try {
targetClass = Class.forName(bucketKeyPair.getBucket().toString());
} catch (Throwable ignored) {
targetClass = Map.class;
}
return (V) getWithMetaData(bucketKeyPair, targetClass).get();
}
public <K> byte[] getAsBytes(K key) {
return getAsBytesWithMetaData(key).get();
}
public <K> RiakValue<byte[]> getAsBytesWithMetaData(K key) {
BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, null);
final RestTemplate restTemplate = getRestTemplate();
if (log.isDebugEnabled()) {
log.debug(String.format("GET object: bucket=%s, key=%s, type=byte[]",
bucketKeyPair.getBucket(),
bucketKeyPair.getKey()));
}
try {
RiakValue<byte[]> bytes = (RiakValue<byte[]>) restTemplate.execute(defaultUri,
HttpMethod.GET,
new RequestCallback() {
public void doWithRequest(ClientHttpRequest request) throws IOException {
List<MediaType> mediaTypes = new ArrayList<MediaType>();
mediaTypes.add(MediaType.APPLICATION_JSON);
request.getHeaders().setAccept(mediaTypes);
}
},
new ResponseExtractor<Object>() {
public Object extractData(ClientHttpResponse response) throws IOException {
InputStream in = response.getBody();
ByteArrayOutputStream out = new ByteArrayOutputStream();
byte[] buff = new byte[in.available()];
for (int bytesRead = in.read(buff); bytesRead > 0; bytesRead = in.read(buff)) {
out.write(buff, 0, bytesRead);
}
HttpHeaders headers = response.getHeaders();
RiakMetaData meta = extractMetaData(headers);
RiakValue<byte[]> val = new RiakValue<byte[]>(out.toByteArray(), meta);
return val;
}
},
bucketKeyPair.getBucket(),
bucketKeyPair.getKey());
if (useCache) {
cache.put(bucketKeyPair, bytes);
}
return bytes;
} catch (HttpClientErrorException e) {
if (e.getStatusCode() != HttpStatus.NOT_FOUND) {
throw new DataStoreOperationException(e.getMessage(), e);
}
}
return null;
}
public <K, T> T getAsType(K key, Class<T> requiredType) {
if (useCache) {
Object obj = checkCache(key, requiredType);
if (null != obj) {
return (T) obj;
}
}
return getWithMetaData(key, requiredType).get();
}
public <K, V> V getAndSet(K key, V value) {
@@ -196,7 +293,7 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
}
public <K> byte[] getAndSetAsBytes(K key, byte[] value) {
byte[] old = getAsType(key, byte[].class);
byte[] old = getAsBytes(key);
setAsBytes(key, value);
return old;
}
@@ -234,6 +331,8 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
return getValuesAsType(keyList, requiredType);
}
/*----------------- Only-Set-Once Operations -----------------*/
public <K, V> KeyValueStoreOperations setIfKeyNonExistent(K key, V value) {
if (!containsKey(key)) {
set(key, value);
@@ -256,6 +355,8 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
return this;
}
/*----------------- Multiple Item Operations -----------------*/
public <K, V> KeyValueStoreOperations setMultiple(Map<K, V> keysAndValues) {
for (Map.Entry<K, V> entry : keysAndValues.entrySet()) {
set(entry.getKey(), entry.getValue());
@@ -284,6 +385,8 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
return this;
}
/*----------------- Key Operations -----------------*/
public <K> boolean containsKey(K key) {
BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, null);
RestTemplate restTemplate = getRestTemplate();
@@ -337,6 +440,52 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
return queue.submit(job);
}
/*----------------- Link Operations -----------------*/
public <K1, K2> RiakTemplate link(K1 destination, K2 source, String tag) {
BucketKeyPair bkpFrom = resolveBucketKeyPair(source, null);
BucketKeyPair bkpTo = resolveBucketKeyPair(destination, null);
RestTemplate restTemplate = getRestTemplate();
RiakValue<byte[]> fromObj = getAsBytesWithMetaData(source);
HttpHeaders headers = new HttpHeaders();
headers.setContentType(fromObj.getMetaData().getContentType());
Object linksObj = fromObj.getMetaData().getProperties().get("Link");
List<String> links = new ArrayList<String>();
if (linksObj instanceof List) {
links.addAll((List) linksObj);
} else if (linksObj instanceof String) {
links.add(linksObj.toString());
}
links.add(String.format("<%s/%s/%s>; riaktag=\"%s\"", extractPrefix(), bkpTo.getBucket(), bkpTo.getKey(), tag));
for (String link : links) {
headers.set("Link", link);
}
HttpEntity entity = new HttpEntity(fromObj.get(), headers);
restTemplate.put(defaultUri, entity, bkpFrom.getBucket(), bkpFrom.getKey());
return this;
}
/*----------------- Bucket Operations -----------------*/
public <B> Map<String, Object> getBucketSchema(B bucket) {
return getBucketSchema(bucket, false);
}
public <B> Map<String, Object> getBucketSchema(B bucket, boolean listKeys) {
RestTemplate restTemplate = getRestTemplate();
ResponseEntity<Map> resp = restTemplate.getForEntity(defaultUri,
Map.class,
bucket,
(listKeys ? "?keys=true" : ""));
if (resp.hasBody()) {
return resp.getBody();
} else {
throw new DataStoreOperationException("Error encountered retrieving bucket schema (Status: " + resp.getStatusCode() + ")");
}
}
public void afterPropertiesSet() throws Exception {
Assert.notNull(conversionService, "Must specify a valid ConversionService.");
if (null == bucketKeyResolvers) {
@@ -346,19 +495,22 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
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) {
ObjectMapper mapper = new ObjectMapper();
CustomSerializerFactory fac = new CustomSerializerFactory();
fac.addSpecificMapping(GStringImpl.class, ToStringSerializer.instance);
mapper.setSerializerFactory(fac);
((MappingJacksonHttpMessageConverter) converter).setObjectMapper(mapper);
}
}
}
}
/*----------------- Utilities -----------------*/
protected BucketKeyPair resolveBucketKeyPair(Object key, Object val) {
BucketKeyResolver resolver = null;
for (BucketKeyResolver r : bucketKeyResolvers) {
@@ -396,4 +548,67 @@ public class RiakTemplate extends RestGatewaySupport implements KeyValueStoreOpe
return mediaType;
}
protected RiakMetaData extractMetaData(HttpHeaders headers) throws IOException {
Map<String, Object> props = new LinkedHashMap<String, Object>();
for (Map.Entry<String, List<String>> entry : headers.entrySet()) {
List<String> val = entry.getValue();
Object prop = (1 == val.size() ? val.get(0) : val);
try {
if (entry.getKey().equals("Last-Modified") || entry.getKey().equals("Date")) {
prop = httpDate.parse(val.get(0));
}
} catch (ParseException e) {
log.error(e.getMessage(), e);
}
if (entry.getKey().equals("Link")) {
List<String> links = new ArrayList<String>();
for (String link : entry.getValue()) {
String[] parts = link.split(",");
for (String part : parts) {
String s = part.replaceAll("<(.+)>; rel=\"(\\S+)\"[,]?", "").trim();
if (!"".equals(s)) {
links.add(s);
}
}
}
props.put("Link", links);
} else {
props.put(entry.getKey().toString(), prop);
}
}
props.put("ETag", headers.getETag());
RiakMetaData meta = new RiakMetaData(headers.getContentType(), props);
return meta;
}
protected <K, T> T checkCache(K key, Class<T> requiredType) {
BucketKeyPair bucketKeyPair = resolveBucketKeyPair(key, requiredType);
RiakValue<?> obj = cache.get(bucketKeyPair);
if (null != obj) {
String bucketName = (null != bucketKeyPair.getBucket() ? bucketKeyPair.getBucket()
.toString() : requiredType.getName());
RestTemplate restTemplate = getRestTemplate();
HttpHeaders resp = restTemplate.headForHeaders(defaultUri, bucketName, bucketKeyPair.getKey());
if (!obj.getMetaData().getProperties().get("ETag").toString().equals(resp.getETag())) {
obj = null;
} else {
if (log.isDebugEnabled()) {
log.debug("Returning CACHED object: " + obj);
}
}
}
return (null != obj ? (T) obj.get() : null);
}
public String extractPrefix() {
Matcher m = prefix.matcher(defaultUri);
if (m.matches()) {
return "/" + m.group(3);
}
return "/riak";
}
}

View File

@@ -0,0 +1,24 @@
package org.springframework.datastore.riak.core;
/**
* @author J. Brisbin <jon@jbrisbin.com>
*/
@SuppressWarnings({"unchecked"})
public class RiakValue<T> implements KeyValueStoreValue {
private Object delegate;
private KeyValueStoreMetaData metaData;
public RiakValue(T delegate, KeyValueStoreMetaData metaData) {
this.delegate = delegate;
this.metaData = metaData;
}
public KeyValueStoreMetaData getMetaData() {
return this.metaData;
}
public T get() {
return (T) delegate;
}
}

View File

@@ -4,7 +4,7 @@ package org.springframework.datastore.riak.core;
* @author J. Brisbin <jon@jbrisbin.com>
*/
@SuppressWarnings({"unchecked"})
public class SimpleBucketKeyPair<B, K> implements BucketKeyPair {
public class SimpleBucketKeyPair<B, K> implements BucketKeyPair, Comparable {
private Object bucket;
private Object key;
@@ -21,4 +21,14 @@ public class SimpleBucketKeyPair<B, K> implements BucketKeyPair {
public K getKey() {
return (K) key;
}
public int compareTo(Object o) {
if (o instanceof SimpleBucketKeyPair) {
SimpleBucketKeyPair pair = (SimpleBucketKeyPair) o;
if (pair.getBucket().equals(bucket) && pair.getKey().equals(key)) {
return 0;
}
}
return -1;
}
}

View File

@@ -64,6 +64,26 @@ class RiakTemplateSpec extends Specification {
}
def "Test getting bucket schema"() {
when:
def schema = riak.getBucketSchema("test", true)
then:
"test" == schema.props.name
}
def "Test get with metadata"() {
when:
def val = riak.getWithMetaData([bucket: "test", key: "test"], LinkedHashMap)
then:
val.metaData.properties["Server"].contains("WebMachine")
}
def "Test containsKey"() {
when:
@@ -74,6 +94,20 @@ class RiakTemplateSpec extends Specification {
}
def "Test linking"() {
given:
riak.link("${TestObject.name}:test", "test:test", "test")
when:
def val = riak.getWithMetaData("test:test", Map)
def result = val.metaData.properties["Link"].collect { it.contains("riaktag=\"test\"") }
then:
1 == result.size()
}
def "Test multiple get"() {
when:

View File

@@ -21,5 +21,6 @@ Import-Template:
org.aopalliance.*;version="[1.0.0, 2.0.0)";resolution:=optional,
org.slf4j.*;version="[1.5.10, 2.0.0)",
org.w3c.dom.*;version="0",
org.codehaus.jackson.*;version="[1.5.6, 1.5.6)",
org.codehaus.jackson.map.*;version="[1.5.6, 1.5.6)",
org.codehaus.groovy.runtime.*;version="[1.7.5, 2.0.0)",