diff --git a/spring-data-riak/README.md b/spring-data-riak/README.md index 4af648c8f..d41439833 100644 --- a/spring-data-riak/README.md +++ b/spring-data-riak/README.md @@ -16,14 +16,10 @@ One cool new feature just added is a Groovy DSL for data access using SDKV/Riak: def result = null riak.set(bucket: "test", key: "test", qos: [dw: "all"], value: obj, wait: 3000L) { + completed(when: { it.integer == 12 }) { result = it.test } + completed { result = "otherwise" } - completed(when: { v -> v.integer == 12 }) { v, meta -> - result = v.test - } - completed { v -> result = "otherwise" } - - failed { e -> println "failure: $e" } - + failed { it.printStackTrace() } } The Groovy DSL will respond to the following methods: @@ -33,17 +29,36 @@ The Groovy DSL will respond to the following methods: * put * get * getAsBytes +* getAsType * containsKey * delete * each -You can nest them, of course. To delete all keys from a bucket using the DSL: +Each completed or failed closure can be accompanied by a "guard" closure. For example, +to process an entry differently, based on the type: - riak.each(bucket: "test") { - completed { v, meta -> - delete(bucket: meta.bucket, key: meta.key) + riak.get(bucket: "test", key: "test") { + completed(when: { it instanceof Map }) { processMap(it) } + completed(when: { it instanceof String }) { processString(it) } + completed(when: { it instanceof byte[] }) { processBytes(it) } + completed { result = "otherwise" } + + failed { it.printStackTrace() } + } + +You can nest them, of course. To insert data and then delete all keys from a bucket: + + riak { + put(bucket: "test", value: [test: "value 1"]) + put(bucket: "test", value: [test: "value 2"]) + put(bucket: "test", value: [test: "value 3"]) + + each(bucket: "test") { + completed { v, meta -> + delete(bucket: meta.bucket, key: meta.key) + } + failed { it.printStackTrace() } } - failed { e -> println "failure: $e" } } 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 6a37f0449..d3acd9f18 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 @@ -23,6 +23,7 @@ import groovy.util.BuilderSupport; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.data.keyvalue.riak.DataStoreOperationException; import org.springframework.data.keyvalue.riak.core.AsyncRiakTemplate; import org.springframework.data.keyvalue.riak.core.RiakQosParameters; @@ -98,6 +99,20 @@ public class RiakBuilder extends BuilderSupport { op.setKey((null != o ? o.toString() : null)); o = attributes.get("value"); op.setValue(o); + o = attributes.get("type"); + if (null != o) { + if (o instanceof Class) { + op.setRequiredType((Class) o); + } else if (o instanceof String) { + try { + op.setRequiredType(Class.forName((String) o)); + } catch (ClassNotFoundException e) { + throw new DataStoreOperationException(e.getMessage(), e); + } + } else { + op.setRequiredType(o.getClass()); + } + } o = attributes.get("qos"); if (null != o) { RiakQosParameters qos = new RiakQosParameters(); @@ -180,7 +195,7 @@ public class RiakBuilder extends BuilderSupport { @Override protected void nodeCompleted(Object parent, Object node) { log.debug("nodeCompleted: " + parent + " " + node); - if (null == parent && node instanceof RiakOperation) { + if (node instanceof RiakOperation) { RiakOperation op = (RiakOperation) node; try { op.call(); 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 57be0a33f..0a22420eb 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 @@ -39,7 +39,7 @@ import java.util.concurrent.*; public class RiakOperation implements Callable { static enum Type { - SET, SETASBYTES, PUT, GET, GETASBYTES, CONTAINSKEY, DELETE, EACH + SET, SETASBYTES, PUT, GET, GETASBYTES, GETASTYPE, CONTAINSKEY, DELETE, EACH } static String COMPLETED = "completed"; @@ -52,6 +52,7 @@ public class RiakOperation implements Callable { protected String bucket; protected String key; protected T value; + protected Class requiredType = null; protected long timeout = -1L; protected QosParameters qosParameters; protected Map> callbacks = new LinkedHashMap>(); @@ -102,6 +103,14 @@ public class RiakOperation implements Callable { this.value = value; } + public Class getRequiredType() { + return requiredType; + } + + public void setRequiredType(Class requiredType) { + this.requiredType = requiredType; + } + public long getTimeout() { return timeout; } @@ -129,6 +138,9 @@ public class RiakOperation implements Callable { case GETASBYTES: f = riak.getAsBytes(bucket, key, callbackInvoker); break; + case GETASTYPE: + f = riak.getAsType(bucket, key, requiredType, callbackInvoker); + break; case PUT: f = riak.put(bucket, value, callbackInvoker); break; 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 10c34d6bb..93ba07b91 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 @@ -43,13 +43,12 @@ class RiakBuilderSpec extends Specification { def result = null when: - riak.set(bucket: "test", key: "test", qos: [dw: "all"], value: obj) { - - completed(when: { it.integer == 12 }) { result = it.test } - completed { result = "otherwise" } - - failed { it.printStackTrace() } - + riak { + set(bucket: "test", key: "test", qos: [dw: "all"], value: obj) { + completed(when: { it.integer == 12 }) { result = it.test } + completed { result = "otherwise" } + failed { it.printStackTrace() } + } } then: @@ -65,12 +64,27 @@ class RiakBuilderSpec extends Specification { when: riak.get(bucket: "test", key: "test") { - completed(when: { it.integer == 12 }) { result = it.test } completed { result = "otherwise" } - failed { it.printStackTrace() } + } + then: + "value" == result + + } + + def "Test builder getAsType"() { + + given: + def riak = new RiakBuilder(riakTemplate) + def result = null + + when: + riak.getAsType(bucket: "test", key: "test", type: Map) { + completed(when: { it instanceof Map }) { result = it.test } + completed { result = "otherwise" } + failed { it.printStackTrace() } } then: @@ -150,6 +164,30 @@ class RiakBuilderSpec extends Specification { } + def "Test builder batch operations"() { + + given: + def riak = new RiakBuilder(riakTemplate) + def ids = [] + + when: + riak { + put(bucket: "test", value: [test: "value 1"]) + put(bucket: "test", value: [test: "value 2"]) + put(bucket: "test", value: [test: "value 3"]) + + each(bucket: "test") { + completed { v, meta -> ids << meta.key } + failed { it.printStackTrace() } + } + } + + then: + null != ids + 3 <= ids.size() + + } + def "Test builder delete"() { given: