diff --git a/spring-data-riak/README.md b/spring-data-riak/README.md index 9edf47f1d..33f5303de 100644 --- a/spring-data-riak/README.md +++ b/spring-data-riak/README.md @@ -33,7 +33,7 @@ The Groovy DSL will respond to the following methods: * getAsType * containsKey * delete -* each +* foreach Each completed or failed closure can be accompanied by a "guard" closure. For example, to process an entry differently, based on the type: diff --git a/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/AsyncBucketKeyValueStoreOperations.java b/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/AsyncBucketKeyValueStoreOperations.java index b006fea00..0b222d855 100644 --- a/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/AsyncBucketKeyValueStoreOperations.java +++ b/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/AsyncBucketKeyValueStoreOperations.java @@ -37,7 +37,7 @@ public interface AsyncBucketKeyValueStoreOperations { * @param value * @param callback Called with the update value pulled from Riak */ - Future set(B bucket, K key, V value, AsyncKeyValueStoreOperation callback); + Future set(B bucket, K key, V value, AsyncKeyValueStoreOperation callback); /** * @param bucket @@ -46,7 +46,7 @@ public interface AsyncBucketKeyValueStoreOperations { * @param qosParams * @return */ - Future set(B bucket, K key, V value, QosParameters qosParams, AsyncKeyValueStoreOperation callback); + Future set(B bucket, K key, V value, QosParameters qosParams, AsyncKeyValueStoreOperation callback); /** * @param bucket @@ -54,7 +54,7 @@ public interface AsyncBucketKeyValueStoreOperations { * @param value * @return */ - Future setAsBytes(B bucket, K key, byte[] value, AsyncKeyValueStoreOperation callback); + Future setAsBytes(B bucket, K key, byte[] value, AsyncKeyValueStoreOperation callback); /** * @param bucket @@ -63,7 +63,7 @@ public interface AsyncBucketKeyValueStoreOperations { * @param qosParams * @return */ - Future setAsBytes(B bucket, K key, byte[] value, QosParameters qosParams, AsyncKeyValueStoreOperation callback); + Future setAsBytes(B bucket, K key, byte[] value, QosParameters qosParams, AsyncKeyValueStoreOperation callback); /** * @param bucket @@ -72,7 +72,7 @@ public interface AsyncBucketKeyValueStoreOperations { * @param metaData * @return */ - Future setWithMetaData(B bucket, K key, V value, Map metaData, AsyncKeyValueStoreOperation callback); + Future setWithMetaData(B bucket, K key, V value, Map metaData, AsyncKeyValueStoreOperation callback); /** * @param bucket @@ -82,21 +82,21 @@ public interface AsyncBucketKeyValueStoreOperations { * @param qosParams * @return */ - Future setWithMetaData(B bucket, K key, V value, Map metaData, QosParameters qosParams, AsyncKeyValueStoreOperation callback); + Future setWithMetaData(B bucket, K key, V value, Map metaData, QosParameters qosParams, AsyncKeyValueStoreOperation callback); /** * @param bucket * @param key * @return */ - Future get(B bucket, K key, AsyncKeyValueStoreOperation callback); + Future get(B bucket, K key, AsyncKeyValueStoreOperation callback); /** * @param bucket * @param key * @return */ - Future getAsBytes(B bucket, K key, AsyncKeyValueStoreOperation callback); + Future getAsBytes(B bucket, K key, AsyncKeyValueStoreOperation callback); /** * @param bucket @@ -104,7 +104,7 @@ public interface AsyncBucketKeyValueStoreOperations { * @param requiredType * @return */ - Future getAsType(B bucket, K key, Class requiredType, AsyncKeyValueStoreOperation callback); + Future getAsType(B bucket, K key, Class requiredType, AsyncKeyValueStoreOperation callback); /** * @param bucket @@ -112,7 +112,7 @@ public interface AsyncBucketKeyValueStoreOperations { * @param value * @return */ - Future getAndSet(B bucket, K key, V value, AsyncKeyValueStoreOperation callback); + Future getAndSet(B bucket, K key, V value, AsyncKeyValueStoreOperation callback); /** * @param bucket @@ -120,7 +120,7 @@ public interface AsyncBucketKeyValueStoreOperations { * @param value * @return */ - Future getAndSetAsBytes(B bucket, K key, byte[] value, AsyncKeyValueStoreOperation callback); + Future getAndSetAsBytes(B bucket, K key, byte[] value, AsyncKeyValueStoreOperation callback); /** * @param bucket @@ -129,7 +129,7 @@ public interface AsyncBucketKeyValueStoreOperations { * @param requiredType * @return */ - Future getAndSetAsType(B bucket, K key, V value, Class requiredType, AsyncKeyValueStoreOperation callback); + Future getAndSetAsType(B bucket, K key, V value, Class requiredType, AsyncKeyValueStoreOperation callback); /** * @param bucket @@ -137,7 +137,7 @@ public interface AsyncBucketKeyValueStoreOperations { * @param value * @return */ - Future setIfKeyNonExistent(B bucket, K key, V value, AsyncKeyValueStoreOperation callback); + Future setIfKeyNonExistent(B bucket, K key, V value, AsyncKeyValueStoreOperation callback); /** * @param bucket @@ -145,14 +145,14 @@ public interface AsyncBucketKeyValueStoreOperations { * @param value * @return */ - Future setIfKeyNonExistentAsBytes(B bucket, K key, byte[] value, AsyncKeyValueStoreOperation callback); + Future setIfKeyNonExistentAsBytes(B bucket, K key, byte[] value, AsyncKeyValueStoreOperation callback); /** * @param bucket * @param key * @return */ - Future containsKey(B bucket, K key, AsyncKeyValueStoreOperation callback); + Future containsKey(B bucket, K key, AsyncKeyValueStoreOperation callback); /** * Delete a specific entry from this data store. @@ -161,6 +161,6 @@ public interface AsyncBucketKeyValueStoreOperations { * @param key * @return */ - Future delete(B bucket, K key, AsyncKeyValueStoreOperation callback); + Future delete(B bucket, K key, AsyncKeyValueStoreOperation callback); } diff --git a/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/AsyncKeyValueStoreOperation.java b/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/AsyncKeyValueStoreOperation.java index cdd6e41fc..291d8d126 100644 --- a/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/AsyncKeyValueStoreOperation.java +++ b/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/AsyncKeyValueStoreOperation.java @@ -21,9 +21,9 @@ package org.springframework.data.keyvalue.riak.core; /** * @author J. Brisbin */ -public interface AsyncKeyValueStoreOperation { +public interface AsyncKeyValueStoreOperation { - void completed(KeyValueStoreMetaData meta, V result); + T completed(KeyValueStoreMetaData meta, V result); - void failed(Throwable error); + T failed(Throwable error); } diff --git a/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/AsyncRiakTemplate.java b/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/AsyncRiakTemplate.java index 09145e747..c3f5c3a66 100644 --- a/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/AsyncRiakTemplate.java +++ b/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/core/AsyncRiakTemplate.java @@ -38,10 +38,7 @@ import java.net.URI; import java.util.ArrayList; import java.util.List; import java.util.Map; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.Future; +import java.util.concurrent.*; /** * @author J. Brisbin @@ -51,7 +48,7 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck protected final Logger log = LoggerFactory.getLogger(getClass()); protected ExecutorService workerPool = Executors.newCachedThreadPool(); - protected AsyncKeyValueStoreOperation defaultErrorHandler = new LoggingErrorHandler(); + protected AsyncKeyValueStoreOperation defaultErrorHandler = new LoggingErrorHandler(); public AsyncRiakTemplate() { super(); @@ -69,28 +66,28 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck this.workerPool = workerPool; } - public AsyncKeyValueStoreOperation getDefaultErrorHandler() { + public AsyncKeyValueStoreOperation getDefaultErrorHandler() { return defaultErrorHandler; } - public void setDefaultErrorHandler(AsyncKeyValueStoreOperation defaultErrorHandler) { + public void setDefaultErrorHandler(AsyncKeyValueStoreOperation defaultErrorHandler) { this.defaultErrorHandler = defaultErrorHandler; } - public Future set(B bucket, K key, V value, AsyncKeyValueStoreOperation callback) { + public Future set(B bucket, K key, V value, AsyncKeyValueStoreOperation callback) { return setWithMetaData(bucket, key, value, null, null, callback); } - public Future set(B bucket, K key, V value, QosParameters qosParams, AsyncKeyValueStoreOperation callback) { + public Future set(B bucket, K key, V value, QosParameters qosParams, AsyncKeyValueStoreOperation callback) { return setWithMetaData(bucket, key, value, null, qosParams, callback); } - public Future setAsBytes(B bucket, K key, byte[] value, AsyncKeyValueStoreOperation callback) { + public Future setAsBytes(B bucket, K key, byte[] value, AsyncKeyValueStoreOperation callback) { return setWithMetaData(bucket, key, value, null, null, callback); } @SuppressWarnings({"unchecked"}) - public Future setWithMetaData(B bucket, K key, V value, Map metaData, QosParameters qosParams, AsyncKeyValueStoreOperation callback) { + public Future setWithMetaData(B bucket, K key, V value, Map metaData, QosParameters qosParams, AsyncKeyValueStoreOperation callback) { String bucketName = (null != bucket ? bucket.toString() : value.getClass().getName()); // Get a key name that may or may not include the QOS parameters. Assert.notNull(key, "Cannot use a key."); @@ -110,22 +107,22 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck } headers.set(RIAK_META_CLASSNAME, value.getClass().getName()); HttpEntity entity = new HttpEntity(value, headers); - return (Future) workerPool.submit(new AsyncPost(bucketName, + return (Future) workerPool.submit(new AsyncPost(bucketName, keyName, entity, callback)); } - public Future put(B bucket, V value, AsyncKeyValueStoreOperation callback) { + public Future put(B bucket, V value, AsyncKeyValueStoreOperation callback) { return put(bucket, value, null, null, callback); } - public Future put(B bucket, V value, Map metaData, AsyncKeyValueStoreOperation callback) { + public Future put(B bucket, V value, Map metaData, AsyncKeyValueStoreOperation callback) { return put(bucket, value, metaData, null, callback); } @SuppressWarnings({"unchecked"}) - public Future put(B bucket, V value, Map metaData, QosParameters qosParams, AsyncKeyValueStoreOperation callback) { + public Future put(B bucket, V value, Map metaData, QosParameters qosParams, AsyncKeyValueStoreOperation callback) { Assert.notNull(bucket, "Bucket cannot be null"); String bucketName = (null != qosParams ? bucket.toString() + extractQosParameters(qosParams) : bucket .toString()); @@ -134,10 +131,10 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck headers.setContentType(extractMediaType(value)); headers.set(RIAK_META_CLASSNAME, value.getClass().getName()); HttpEntity entity = new HttpEntity(value, headers); - return (Future) workerPool.submit(new AsyncPut(bucketName, entity, callback)); + return (Future) workerPool.submit(new AsyncPut(bucketName, entity, callback)); } - public Future get(B bucket, K key, AsyncKeyValueStoreOperation callback) { + public Future get(B bucket, K key, AsyncKeyValueStoreOperation callback) { return getWithMetaData(bucket, key, null, callback); } @@ -158,7 +155,7 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck } @SuppressWarnings({"unchecked"}) - public Future getBucketSchema(B bucket, QosParameters qosParams, final AsyncKeyValueStoreOperation> callback) { + public Future getBucketSchema(B bucket, QosParameters qosParams, final AsyncKeyValueStoreOperation, R> callback) { Assert.notNull(bucket, "Bucket cannot be null"); Assert.notNull(callback, "Callback cannot be null"); @@ -168,20 +165,20 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck return workerPool.submit(new AsyncGet(bucketName, "?keys=true", Map.class, - new AsyncKeyValueStoreOperation() { + new AsyncKeyValueStoreOperation() { @SuppressWarnings({"unchecked"}) - public void completed(KeyValueStoreMetaData meta, Object result) { - callback.completed(meta, (Map) result); + public Object completed(KeyValueStoreMetaData meta, Object result) { + return callback.completed(meta, (Map) result); } - public void failed(Throwable error) { - callback.failed(error); + public Object failed(Throwable error) { + return callback.failed(error); } })); } @SuppressWarnings({"unchecked"}) - public Future getWithMetaData(B bucket, K key, Class requiredType, AsyncKeyValueStoreOperation callback) { + public Future getWithMetaData(B bucket, K key, Class requiredType, AsyncKeyValueStoreOperation callback) { String bucketName = (null != bucket ? bucket.toString() : requiredType.getName()); // Get a key name that may or may not include the QOS parameters. Assert.notNull(key, "Cannot use a null key."); @@ -190,32 +187,32 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck if (null == requiredType) { requiredType = (Class) getType(bucketName, key.toString()); } - return workerPool.submit(new AsyncGet(bucketName, + return workerPool.submit(new AsyncGet(bucketName, key.toString(), requiredType, callback)); } - public Future getAsBytes(B bucket, K key, AsyncKeyValueStoreOperation callback) { + public Future getAsBytes(B bucket, K key, AsyncKeyValueStoreOperation callback) { return getWithMetaData(bucket, key, byte[].class, callback); } - public Future getAsType(B bucket, K key, Class requiredType, AsyncKeyValueStoreOperation callback) { + public Future getAsType(B bucket, K key, Class requiredType, AsyncKeyValueStoreOperation callback) { return getWithMetaData(bucket, key, requiredType, callback); } - public Future getAndSet(final B bucket, final K key, final V value, final AsyncKeyValueStoreOperation callback) { + public Future getAndSet(final B bucket, final K key, final V value, final AsyncKeyValueStoreOperation callback) { final List> futures = new ArrayList>(); try { - getWithMetaData(bucket, key, null, new AsyncKeyValueStoreOperation() { + getWithMetaData(bucket, key, null, new AsyncKeyValueStoreOperation() { @SuppressWarnings({"unchecked"}) - public void completed(KeyValueStoreMetaData meta, Object result) { + public Object completed(KeyValueStoreMetaData meta, Object result) { futures.add(setWithMetaData(bucket, key, value, null, null, null)); - callback.completed(meta, (V) result); + return callback.completed(meta, (V) result); } - public void failed(Throwable error) { - callback.failed(error); + public Object failed(Throwable error) { + return callback.failed(error); } }).get(); } catch (InterruptedException e) { @@ -226,88 +223,98 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck return futures.size() > 0 ? futures.get(0) : null; } - public Future getAndSetAsBytes(B bucket, K key, byte[] value, AsyncKeyValueStoreOperation callback) { + public Future getAndSetAsBytes(B bucket, K key, byte[] value, AsyncKeyValueStoreOperation callback) { return getAndSet(bucket, key, value, callback); } - public Future getAndSetAsType(final B bucket, final K key, final V value, final Class requiredType, final AsyncKeyValueStoreOperation callback) { + public Future getAndSetAsType(final B bucket, final K key, final V value, final Class requiredType, final AsyncKeyValueStoreOperation callback) { final List> futures = new ArrayList>(); - getWithMetaData(bucket, key, requiredType, new AsyncKeyValueStoreOperation() { + getWithMetaData(bucket, key, requiredType, new AsyncKeyValueStoreOperation() { @SuppressWarnings({"unchecked"}) - public void completed(KeyValueStoreMetaData meta, T result) { - futures.add(setWithMetaData(bucket, key, value, null, null, null)); - callback.completed(meta, (V) result); + public R completed(KeyValueStoreMetaData meta, T result) { + try { + setWithMetaData(bucket, key, value, null, null, null).get(); + return callback.completed(meta, result); + } catch (InterruptedException e) { + return callback.failed(e); + } catch (ExecutionException e) { + return callback.failed(e); + } } - public void failed(Throwable error) { - callback.failed(error); + public R failed(Throwable error) { + return callback.failed(error); } }); return futures.size() > 0 ? futures.get(0) : null; } - public Future setIfKeyNonExistent(final B bucket, final K key, final V value, final AsyncKeyValueStoreOperation callback) { - return containsKey(bucket, key, new AsyncKeyValueStoreOperation() { - public void completed(KeyValueStoreMetaData meta, Boolean result) { + public Future setIfKeyNonExistent(final B bucket, final K key, final V value, final AsyncKeyValueStoreOperation callback) { + return containsKey(bucket, key, new AsyncKeyValueStoreOperation() { + public Object completed(KeyValueStoreMetaData meta, Boolean result) { if (!result) { - setWithMetaData(bucket, key, value, null, null, callback); + return setWithMetaData(bucket, key, value, null, null, callback); + } else { + return null; } } - public void failed(Throwable error) { - callback.failed(error); + public Object failed(Throwable error) { + return callback.failed(error); } }); } - public Future setIfKeyNonExistentAsBytes(final B bucket, final K key, final byte[] value, final AsyncKeyValueStoreOperation callback) { - return containsKey(bucket, key, new AsyncKeyValueStoreOperation() { - public void completed(KeyValueStoreMetaData meta, Boolean result) { + public Future setIfKeyNonExistentAsBytes(final B bucket, final K key, final byte[] value, final AsyncKeyValueStoreOperation callback) { + return containsKey(bucket, key, new AsyncKeyValueStoreOperation() { + public Object completed(KeyValueStoreMetaData meta, Boolean result) { if (!result) { - setWithMetaData(bucket, key, value, null, null, callback); + return setWithMetaData(bucket, key, value, null, null, callback); + } else { + return null; } } - public void failed(Throwable error) { - callback.failed(error); + public Object failed(Throwable error) { + return callback.failed(error); } }); } - public Future containsKey(B bucket, K key, final AsyncKeyValueStoreOperation callback) { + public Future containsKey(B bucket, K key, final AsyncKeyValueStoreOperation callback) { Assert.notNull(bucket, "Bucket cannot be null when checking for existence."); Assert.notNull(key, "Key cannot be null when checking for existence"); return workerPool.submit(new AsyncHead(bucket.toString(), key.toString(), - new AsyncKeyValueStoreOperation() { - public void completed(KeyValueStoreMetaData meta, HttpHeaders result) { - callback.completed(null, (null != result)); + new AsyncKeyValueStoreOperation() { + public Object completed(KeyValueStoreMetaData meta, HttpHeaders result) { + return callback.completed(null, (null != result)); } - public void failed(Throwable error) { - callback.failed(error); + public Object failed(Throwable error) { + return callback.failed(error); } })); } - public Future delete(B bucket, K key, AsyncKeyValueStoreOperation callback) { + public Future delete(B bucket, K key, AsyncKeyValueStoreOperation callback) { Assert.notNull(bucket, "Bucket cannot be null when deleting."); Assert.notNull(key, "Key cannot be null when deleting."); return workerPool.submit(new AsyncDelete(bucket.toString(), key.toString(), callback)); } - public Future setAsBytes(B bucket, K key, byte[] value, QosParameters qosParams, AsyncKeyValueStoreOperation callback) { + public Future setAsBytes(B bucket, K key, byte[] value, QosParameters qosParams, AsyncKeyValueStoreOperation callback) { return setWithMetaData(bucket, key, value, null, qosParams, callback); } - public Future setWithMetaData(B bucket, K key, V value, Map metaData, AsyncKeyValueStoreOperation callback) { + public Future setWithMetaData(B bucket, K key, V value, Map metaData, AsyncKeyValueStoreOperation callback) { return setWithMetaData(bucket, key, value, metaData, null, callback); } /* ---------------- Map/Reduce ---------------- */ @SuppressWarnings({"unchecked"}) - public Future execute(MapReduceJob job, AsyncKeyValueStoreOperation> callback) { + public Future execute(MapReduceJob job, AsyncKeyValueStoreOperation, R> callback) { HttpHeaders headers = defaultHeaders(null); headers.setContentType(MediaType.APPLICATION_JSON); HttpEntity json = new HttpEntity(job.toJson(), headers); @@ -315,19 +322,19 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck } /* ---------------- Runnable helpers ---------------- */ - protected class AsyncPut implements Runnable { + protected class AsyncPut implements Callable { private String bucket; private HttpEntity entity = null; - private AsyncKeyValueStoreOperation callback = null; + private AsyncKeyValueStoreOperation callback = null; - public AsyncPut(String bucket, HttpEntity entity, AsyncKeyValueStoreOperation callback) { + public AsyncPut(String bucket, HttpEntity entity, AsyncKeyValueStoreOperation callback) { this.bucket = bucket; this.entity = entity; this.callback = callback; } - public void run() { + public R call() throws Exception { try { URI location = getRestTemplate().postForLocation(defaultUri, entity, bucket, ""); String path = location.getPath(); @@ -338,28 +345,29 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck RiakMetaData meta = extractMetaData(headers); meta.setBucket((null != bucket ? bucket.toString() : null)); meta.setKey((null != key ? key.toString() : null)); - callback.completed(meta, entity.getBody()); + return callback.completed(meta, entity.getBody()); } } catch (Throwable t) { DataStoreOperationException dsoe = new DataStoreOperationException(t.getMessage(), t); if (null != callback) { - callback.failed(dsoe); + return callback.failed(dsoe); } else { defaultErrorHandler.failed(dsoe); } } + return null; } } - protected class AsyncPost implements Runnable { + protected class AsyncPost implements Callable { private String bucket; private String key; private HttpEntity entity = null; - private AsyncKeyValueStoreOperation callback = null; + private AsyncKeyValueStoreOperation callback = null; - public AsyncPost(String bucket, String key, HttpEntity entity, AsyncKeyValueStoreOperation callback) { + public AsyncPost(String bucket, String key, HttpEntity entity, AsyncKeyValueStoreOperation callback) { this.bucket = bucket; this.key = key; this.entity = entity; @@ -367,7 +375,7 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck } @SuppressWarnings({"unchecked"}) - public void run() { + public R call() throws Exception { try { HttpEntity result = getRestTemplate().postForEntity(defaultUri, entity, @@ -384,32 +392,33 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck RiakMetaData meta = extractMetaData(result.getHeaders()); meta.setBucket((null != bucket ? bucket.toString() : null)); meta.setKey((null != key ? key.toString() : null)); - callback.completed(meta, (V) result.getBody()); + return callback.completed(meta, (V) result.getBody()); } } catch (Throwable t) { DataStoreOperationException dsoe = new DataStoreOperationException(t.getMessage(), t); if (null != callback) { - callback.failed(dsoe); + return callback.failed(dsoe); } else { defaultErrorHandler.failed(dsoe); } } + return null; } } - protected class AsyncMapReduce implements Runnable { + protected class AsyncMapReduce implements Callable { private HttpEntity entity = null; - private AsyncKeyValueStoreOperation> callback = null; + private AsyncKeyValueStoreOperation, R> callback = null; - public AsyncMapReduce(HttpEntity entity, AsyncKeyValueStoreOperation> callback) { + public AsyncMapReduce(HttpEntity entity, AsyncKeyValueStoreOperation, R> callback) { this.entity = entity; this.callback = callback; } @SuppressWarnings({"unchecked"}) - public void run() { + public R call() throws Exception { try { HttpEntity result = getRestTemplate().postForEntity(mapReduceUri, entity, @@ -419,35 +428,36 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck } if (null != callback) { RiakMetaData meta = extractMetaData(result.getHeaders()); - callback.completed(meta, result.getBody()); + return callback.completed(meta, result.getBody()); } } catch (Throwable t) { DataStoreOperationException dsoe = new DataStoreOperationException(t.getMessage(), t); if (null != callback) { - callback.failed(dsoe); + return callback.failed(dsoe); } else { defaultErrorHandler.failed(dsoe); } } + return null; } } - protected class AsyncGet implements Runnable { + protected class AsyncGet implements Callable { private String bucket; private String key; private Class requiredType; - private AsyncKeyValueStoreOperation callback = null; + private AsyncKeyValueStoreOperation callback = null; - public AsyncGet(String bucket, String key, Class requiredType, AsyncKeyValueStoreOperation callback) { + public AsyncGet(String bucket, String key, Class requiredType, AsyncKeyValueStoreOperation callback) { this.bucket = bucket; this.key = key; this.requiredType = requiredType; this.callback = callback; } - public void run() { + public R call() throws Exception { try { ResponseEntity result = getRestTemplate().getForEntity(defaultUri, requiredType, @@ -462,7 +472,7 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck cache.put(new SimpleBucketKeyPair(bucket, key), val); } if (null != callback) { - callback.completed(meta, val.get()); + return callback.completed(meta, val.get()); } if (log.isDebugEnabled()) { log.debug(String.format("GET object: bucket=%s, key=%s, type=%s", @@ -474,80 +484,85 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck } catch (Throwable t) { DataStoreOperationException dsoe = new DataStoreOperationException(t.getMessage(), t); if (null != callback) { - callback.failed(dsoe); + return callback.failed(dsoe); } else { defaultErrorHandler.failed(dsoe); } } + return null; } } - protected class AsyncHead implements Runnable { + protected class AsyncHead implements Callable { private String bucket; private String key; - private AsyncKeyValueStoreOperation callback = null; + private AsyncKeyValueStoreOperation callback = null; - public AsyncHead(String bucket, String key, AsyncKeyValueStoreOperation callback) { + public AsyncHead(String bucket, String key, AsyncKeyValueStoreOperation callback) { this.bucket = bucket; this.key = key; this.callback = callback; } - public void run() { + public R call() throws Exception { try { HttpHeaders headers = getRestTemplate().headForHeaders(defaultUri, bucket, key); if (null != headers) { if (null != callback) { - callback.completed(null, headers); + return callback.completed(null, headers); } } } catch (Throwable t) { DataStoreOperationException dsoe = new DataStoreOperationException(t.getMessage(), t); if (null != callback) { - callback.failed(dsoe); + return callback.failed(dsoe); } else { defaultErrorHandler.failed(dsoe); } } + return null; } } - protected class AsyncDelete implements Runnable { + protected class AsyncDelete implements Callable { private String bucket; private String key; - private AsyncKeyValueStoreOperation callback = null; + private AsyncKeyValueStoreOperation callback = null; - public AsyncDelete(String bucket, String key, AsyncKeyValueStoreOperation callback) { + public AsyncDelete(String bucket, String key, AsyncKeyValueStoreOperation callback) { this.bucket = bucket; this.key = key; this.callback = callback; } - public void run() { + public R call() throws Exception { try { getRestTemplate().delete(defaultUri, bucket, key); if (null != callback) { - callback.completed(null, true); + return callback.completed(null, true); } } catch (Throwable t) { DataStoreOperationException dsoe = new DataStoreOperationException(t.getMessage(), t); if (null != callback) { - callback.failed(dsoe); + return callback.failed(dsoe); } else { defaultErrorHandler.failed(dsoe); } } + return null; } } - protected class LoggingErrorHandler implements AsyncKeyValueStoreOperation { - public void completed(KeyValueStoreMetaData meta, Throwable result) { + protected class LoggingErrorHandler implements AsyncKeyValueStoreOperation { + public Object completed(KeyValueStoreMetaData meta, Throwable result) { + return null; } - public void failed(Throwable error) { + public Object failed(Throwable error) { log.error(error.getMessage(), error); + return null; } } diff --git a/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/groovy/RiakBuilder.java b/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/groovy/RiakBuilder.java index e5bbb3ce8..2c17e027b 100644 --- a/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/groovy/RiakBuilder.java +++ b/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/groovy/RiakBuilder.java @@ -30,6 +30,7 @@ import org.springframework.data.keyvalue.riak.core.SimpleBucketKeyPair; import org.springframework.data.keyvalue.riak.mapreduce.*; import java.util.ArrayList; +import java.util.LinkedList; import java.util.List; import java.util.Map; import java.util.concurrent.ExecutorService; @@ -46,6 +47,7 @@ public class RiakBuilder extends BuilderSupport { @Autowired(required = false) protected ExecutorService workerPool = Executors.newCachedThreadPool(); protected String defaultBucketName; + protected List results = new LinkedList(); public RiakBuilder() { } @@ -297,7 +299,11 @@ public class RiakBuilder extends BuilderSupport { } return oper; } + } else if ("call".equals(methodName)) { + results.clear(); + defaultBucketName = null; } + // By default return super.invokeMethod(methodName, arg); } @@ -305,14 +311,7 @@ public class RiakBuilder extends BuilderSupport { @Override protected void nodeCompleted(Object parent, Object node) { log.debug("nodeCompleted: parent=" + parent + ", node=" + node); - if (node instanceof RiakOperation) { - RiakOperation op = (RiakOperation) node; - try { - op.call(); - } catch (Exception e) { - log.error(e.getMessage(), e); - } - } else if (parent instanceof RiakMapReduceOperation && node instanceof QueryPhase) { + if (parent instanceof RiakMapReduceOperation && node instanceof QueryPhase) { QueryPhase p = (QueryPhase) node; MapReduceOperation oper = null; if ("javascript".equals(p.language)) { @@ -341,15 +340,28 @@ public class RiakBuilder extends BuilderSupport { @Override protected Object postNodeCompletion(Object parent, Object node) { log.debug("postNodeCompletion: " + parent + " " + node); - if (null == parent && node instanceof RiakMapReduceOperation) { - RiakMapReduceOperation oper = (RiakMapReduceOperation) node; + if (node instanceof RiakOperation) { + RiakOperation op = (RiakOperation) node; try { - return oper.call(); + Object o = op.call(); + if (null != o) { + results.add(o); + } + return o; + } catch (Exception e) { + log.error(e.getMessage(), e); + } + } else if (null == parent && node instanceof RiakMapReduceOperation) { + RiakMapReduceOperation oper = (RiakMapReduceOperation) node; + try { + Object o = oper.call(); + if (null != o) { + results.add(o); + } + return o; } catch (Exception e) { log.error(e.getMessage(), e); } - } else if (null == parent && node == parent) { - defaultBucketName = null; } return super.postNodeCompletion(parent, node); diff --git a/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/groovy/RiakMapReduceOperation.java b/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/groovy/RiakMapReduceOperation.java index 2e72b0b0c..7d7ac3df8 100644 --- a/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/groovy/RiakMapReduceOperation.java +++ b/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/groovy/RiakMapReduceOperation.java @@ -82,17 +82,19 @@ public class RiakMapReduceOperation implements Callable { } public Object call() throws Exception { - Future f = riak.execute(job, new AsyncKeyValueStoreOperation>() { - public void completed(KeyValueStoreMetaData meta, List result) { + Future f = riak.execute(job, new AsyncKeyValueStoreOperation, Object>() { + public Object completed(KeyValueStoreMetaData meta, List result) { Object arg = new Object[]{result, meta}; if (null != completed) { - completed.call(arg); + return completed.call(arg); + } else { + return new Object[]{result, meta}; } } - public void failed(Throwable error) { + public Object failed(Throwable error) { if (null != failed) { - failed.call(error); + return failed.call(error); } else { throw new RuntimeException(error); } diff --git a/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/groovy/RiakOperation.java b/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/groovy/RiakOperation.java index 33d93f65d..acdf21308 100644 --- a/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/groovy/RiakOperation.java +++ b/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/groovy/RiakOperation.java @@ -27,10 +27,7 @@ import org.springframework.data.keyvalue.riak.core.AsyncRiakTemplate; import org.springframework.data.keyvalue.riak.core.KeyValueStoreMetaData; import org.springframework.data.keyvalue.riak.core.QosParameters; -import java.util.ArrayList; -import java.util.LinkedHashMap; -import java.util.List; -import java.util.Map; +import java.util.*; import java.util.concurrent.*; /** @@ -129,7 +126,7 @@ public class RiakOperation implements Callable { } @SuppressWarnings({"unchecked"}) - public T call() throws Exception { + public Object call() throws Exception { Future f = null; switch (type) { case GET: @@ -165,16 +162,25 @@ public class RiakOperation implements Callable { case FOREACH: f = riak.getBucketSchema(bucket, null, - new AsyncKeyValueStoreOperation>() { - public void completed(KeyValueStoreMetaData meta, Map result) { + new AsyncKeyValueStoreOperation, Object>() { + public Object completed(KeyValueStoreMetaData meta, Map result) { + List results = new LinkedList(); List keys = (List) result.get("keys"); for (String key : keys) { try { Future getFut = riak.get(bucket, key, callbackInvoker); if (timeout > 0) { - getFut.get(timeout, TimeUnit.MILLISECONDS); + Object o = getFut.get(timeout, TimeUnit.MILLISECONDS); + if (null != o) { + results.add(o); + } } else if (timeout < 0) { - getFut.get(); + Object o = getFut.get(); + if (null != o) { + results.add(o); + } + } else { + results.add(getFut); } } catch (InterruptedException e) { throw new DataStoreOperationException(e.getMessage(), e); @@ -184,10 +190,11 @@ public class RiakOperation implements Callable { throw new DataStoreOperationException(e.getMessage(), e); } } + return (results.size() > 0 ? results : null); } - public void failed(Throwable error) { - log.error(error.getMessage(), error); + public Object failed(Throwable error) { + throw new RuntimeException(error); } }); break; @@ -196,14 +203,14 @@ public class RiakOperation implements Callable { if (null != f) { if (timeout > 0) { // Block until finished or timeout - return (T) f.get(timeout, TimeUnit.MILLISECONDS); + return f.get(timeout, TimeUnit.MILLISECONDS); } else if (timeout < 0) { // Block indefinitely - return (T) f.get(); + return f.get(); } } - return (T) f; + return f; } class GuardedClosure { @@ -228,9 +235,9 @@ public class RiakOperation implements Callable { class ClosureInvokingCallback implements AsyncKeyValueStoreOperation { - public void completed(KeyValueStoreMetaData meta, Object result) { + public Object completed(KeyValueStoreMetaData meta, Object result) { if (!callbacks.containsKey(COMPLETED)) { - return; + return new Object[]{result, meta}; } for (GuardedClosure cl : callbacks.get(COMPLETED)) { boolean execute = true; @@ -256,21 +263,21 @@ public class RiakOperation implements Callable { if (execute) { Closure callback = cl.getDelegate(); if (callback.getParameterTypes().length == 2) { - // Pass value and metadata - callback.call(new Object[]{result, meta}); + return callback.call(new Object[]{result, meta}); } else { - callback.call(result); + return callback.call(result); } - break; } } + return null; } - public void failed(Throwable error) { + public Object failed(Throwable error) { + if (!callbacks.containsKey(FAILED)) { + throw new RuntimeException(error); + } for (GuardedClosure cl : callbacks.get(FAILED)) { boolean execute = true; - Object param; - Closure guardExpr = cl.getGuard(); if (null != guardExpr) { Object guardResult = guardExpr.call(error); @@ -285,9 +292,10 @@ public class RiakOperation implements Callable { if (execute) { Closure callback = cl.getDelegate(); - callback.call(error); + return callback.call(error); } } + return null; } } diff --git a/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/mapreduce/AsyncMapReduceOperations.java b/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/mapreduce/AsyncMapReduceOperations.java index a42162d90..fffc52c14 100644 --- a/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/mapreduce/AsyncMapReduceOperations.java +++ b/spring-data-riak/src/main/java/org/springframework/data/keyvalue/riak/mapreduce/AsyncMapReduceOperations.java @@ -36,6 +36,6 @@ public interface AsyncMapReduceOperations { * @param job * @return */ - Future execute(MapReduceJob job, AsyncKeyValueStoreOperation> callback); + Future execute(MapReduceJob job, AsyncKeyValueStoreOperation, R> callback); } diff --git a/spring-data-riak/src/test/groovy/org/springframework/data/keyvalue/riak/core/RiakBuilderSpec.groovy b/spring-data-riak/src/test/groovy/org/springframework/data/keyvalue/riak/core/RiakBuilderSpec.groovy index d2e1c9cdf..1b8521aad 100644 --- a/spring-data-riak/src/test/groovy/org/springframework/data/keyvalue/riak/core/RiakBuilderSpec.groovy +++ b/spring-data-riak/src/test/groovy/org/springframework/data/keyvalue/riak/core/RiakBuilderSpec.groovy @@ -18,6 +18,7 @@ package org.springframework.data.keyvalue.riak.core +import java.util.concurrent.Future import org.springframework.beans.factory.annotation.Autowired import org.springframework.context.ApplicationContext import org.springframework.data.keyvalue.riak.groovy.RiakBuilder @@ -70,10 +71,30 @@ class RiakBuilderSpec extends Specification { } then: + null != result "value" == result } + def "Test builder async get"() { + + given: + def riak = new RiakBuilder(riakTemplate) + def result = null + + when: + def f = riak.get(bucket: "test", key: "test", wait: 0) { + completed(when: { it.integer == 12 }) { result = it.test } + completed { result = "otherwise" } + failed { it.printStackTrace() } + } + + then: + f instanceof Future + null != f.get() + + } + def "Test builder getAsType"() { given: @@ -115,17 +136,15 @@ class RiakBuilderSpec extends Specification { given: def riak = new RiakBuilder(riakTemplate) - def result = null when: - riak.get(bucket: "test", key: "test") { - completed { result = it } + def result = riak.get(bucket: "test", key: "test") { failed { it.printStackTrace() } } then: null != result - "test bytes".bytes == result + "test bytes".bytes == result[0] } @@ -134,11 +153,10 @@ class RiakBuilderSpec extends Specification { given: def obj = [test: "value", integer: 12] def riak = new RiakBuilder(riakTemplate) - def id = null when: - riak.put(bucket: "test", qos: [dw: "all"], value: obj) { - completed { v, meta -> id = meta.key } + def id = riak.put(bucket: "test", qos: [dw: "all"], value: obj) { + completed { v, meta -> meta.key } failed { it.printStackTrace() } } @@ -218,7 +236,6 @@ class RiakBuilderSpec extends Specification { given: def riak = new RiakBuilder(riakTemplate) - def result = [] when: riak { @@ -232,14 +249,14 @@ class RiakBuilderSpec extends Specification { source "function(v){ ejsLog('/tmp/mapred.log', JSON.stringify(arguments)); return Riak.reduceSum(v); }" } } - completed { result = it } + completed { it } failed { it.printStackTrace() } } } then: - null != result - 1 <= result.size() + null != riak.results + 1 <= riak.results.size() }