Tweaked README to cover batch updates with Groovy DSL, added getAsType
This commit is contained in:
@@ -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" }
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -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<Object> op = (RiakOperation<Object>) node;
|
||||
try {
|
||||
op.call();
|
||||
|
||||
@@ -39,7 +39,7 @@ import java.util.concurrent.*;
|
||||
public class RiakOperation<T> 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<T> implements Callable {
|
||||
protected String bucket;
|
||||
protected String key;
|
||||
protected T value;
|
||||
protected Class<?> requiredType = null;
|
||||
protected long timeout = -1L;
|
||||
protected QosParameters qosParameters;
|
||||
protected Map<String, List<GuardedClosure>> callbacks = new LinkedHashMap<String, List<GuardedClosure>>();
|
||||
@@ -102,6 +103,14 @@ public class RiakOperation<T> 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<T> 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;
|
||||
|
||||
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user