|
|
|
|
@@ -89,24 +89,31 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck
|
|
|
|
|
return defaultErrorHandler;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public void setDefaultErrorHandler(AsyncKeyValueStoreOperation<Throwable, Object> defaultErrorHandler) {
|
|
|
|
|
public void setDefaultErrorHandler(
|
|
|
|
|
AsyncKeyValueStoreOperation<Throwable, Object> defaultErrorHandler) {
|
|
|
|
|
this.defaultErrorHandler = defaultErrorHandler;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public <B, K, V, R> Future<?> set(B bucket, K key, V value, AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
public <B, K, V, R> Future<?> set(B bucket, K key, V value,
|
|
|
|
|
AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
return setWithMetaData(bucket, key, value, null, null, callback);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public <B, K, V, R> Future<?> set(B bucket, K key, V value, QosParameters qosParams, AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
public <B, K, V, R> Future<?> set(B bucket, K key, V value, QosParameters qosParams,
|
|
|
|
|
AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
return setWithMetaData(bucket, key, value, null, qosParams, callback);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public <B, K, R> Future<?> setAsBytes(B bucket, K key, byte[] value, AsyncKeyValueStoreOperation<byte[], R> callback) {
|
|
|
|
|
public <B, K, R> Future<?> setAsBytes(B bucket, K key, byte[] value,
|
|
|
|
|
AsyncKeyValueStoreOperation<byte[], R> callback) {
|
|
|
|
|
return setWithMetaData(bucket, key, value, null, null, callback);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@SuppressWarnings({"unchecked"})
|
|
|
|
|
public <B, K, V, R> Future<V> setWithMetaData(B bucket, K key, V value, Map<String, String> metaData, QosParameters qosParams, AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
public <B, K, V, R> Future<V> setWithMetaData(B bucket, K key, V value,
|
|
|
|
|
Map<String, String> metaData,
|
|
|
|
|
QosParameters qosParams,
|
|
|
|
|
AsyncKeyValueStoreOperation<V, R> 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 <NULL> key.");
|
|
|
|
|
@@ -116,7 +123,13 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck
|
|
|
|
|
KeyValueStoreMetaData origMeta = getMetaData(bucket, keyName);
|
|
|
|
|
String vclock = null;
|
|
|
|
|
if (null != origMeta) {
|
|
|
|
|
vclock = origMeta.getProperties().get(RIAK_VCLOCK).toString();
|
|
|
|
|
Map<String, Object> mprops = origMeta.getProperties();
|
|
|
|
|
if (null != mprops) {
|
|
|
|
|
Object o = mprops.get(RIAK_VCLOCK);
|
|
|
|
|
if (null != o) {
|
|
|
|
|
vclock = o.toString();
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
HttpHeaders headers = defaultHeaders(metaData);
|
|
|
|
|
@@ -132,18 +145,23 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck
|
|
|
|
|
callback));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public <B, V, R> Future<V> put(B bucket, V value, AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
public <B, V, R> Future<V> put(B bucket, V value,
|
|
|
|
|
AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
return put(bucket, value, null, null, callback);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public <B, V, R> Future<V> put(B bucket, V value, Map<String, String> metaData, AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
public <B, V, R> Future<V> put(B bucket, V value, Map<String, String> metaData,
|
|
|
|
|
AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
return put(bucket, value, metaData, null, callback);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@SuppressWarnings({"unchecked"})
|
|
|
|
|
public <B, V, R> Future<V> put(B bucket, V value, Map<String, String> metaData, QosParameters qosParams, AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
public <B, V, R> Future<V> put(B bucket, V value, Map<String, String> metaData,
|
|
|
|
|
QosParameters qosParams,
|
|
|
|
|
AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
Assert.notNull(bucket, "Bucket cannot be null");
|
|
|
|
|
String bucketName = (null != qosParams ? bucket.toString() + extractQosParameters(qosParams) : bucket
|
|
|
|
|
String bucketName = (null != qosParams ? bucket.toString() + extractQosParameters(
|
|
|
|
|
qosParams) : bucket
|
|
|
|
|
.toString());
|
|
|
|
|
|
|
|
|
|
HttpHeaders headers = defaultHeaders(metaData);
|
|
|
|
|
@@ -153,7 +171,8 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck
|
|
|
|
|
return (Future<V>) workerPool.submit(new AsyncPut<V, R>(bucketName, entity, callback));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public <B, K, V, R> Future<?> get(B bucket, K key, AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
public <B, K, V, R> Future<?> get(B bucket, K key,
|
|
|
|
|
AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
return getWithMetaData(bucket, key, null, callback);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@@ -174,11 +193,13 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@SuppressWarnings({"unchecked"})
|
|
|
|
|
public <B, R> Future<?> getBucketSchema(B bucket, QosParameters qosParams, final AsyncKeyValueStoreOperation<Map<String, Object>, R> callback) {
|
|
|
|
|
public <B, R> Future<?> getBucketSchema(B bucket, QosParameters qosParams,
|
|
|
|
|
final AsyncKeyValueStoreOperation<Map<String, Object>, R> callback) {
|
|
|
|
|
Assert.notNull(bucket, "Bucket cannot be null");
|
|
|
|
|
Assert.notNull(callback, "Callback cannot be null");
|
|
|
|
|
|
|
|
|
|
String bucketName = (null != qosParams ? bucket.toString() + extractQosParameters(qosParams) : bucket
|
|
|
|
|
String bucketName = (null != qosParams ? bucket.toString() + extractQosParameters(
|
|
|
|
|
qosParams) : bucket
|
|
|
|
|
.toString());
|
|
|
|
|
|
|
|
|
|
return workerPool.submit(new AsyncGet(bucketName,
|
|
|
|
|
@@ -197,7 +218,8 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@SuppressWarnings({"unchecked"})
|
|
|
|
|
public <B, K, T, R> Future<?> getWithMetaData(B bucket, K key, Class<T> requiredType, AsyncKeyValueStoreOperation<T, R> callback) {
|
|
|
|
|
public <B, K, T, R> Future<?> getWithMetaData(B bucket, K key, Class<T> requiredType,
|
|
|
|
|
AsyncKeyValueStoreOperation<T, R> 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.");
|
|
|
|
|
@@ -212,15 +234,18 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck
|
|
|
|
|
callback));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public <B, K, R> Future<?> getAsBytes(B bucket, K key, AsyncKeyValueStoreOperation<byte[], R> callback) {
|
|
|
|
|
public <B, K, R> Future<?> getAsBytes(B bucket, K key,
|
|
|
|
|
AsyncKeyValueStoreOperation<byte[], R> callback) {
|
|
|
|
|
return getWithMetaData(bucket, key, byte[].class, callback);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public <B, K, T, R> Future<?> getAsType(B bucket, K key, Class<T> requiredType, AsyncKeyValueStoreOperation<T, R> callback) {
|
|
|
|
|
public <B, K, T, R> Future<?> getAsType(B bucket, K key, Class<T> requiredType,
|
|
|
|
|
AsyncKeyValueStoreOperation<T, R> callback) {
|
|
|
|
|
return getWithMetaData(bucket, key, requiredType, callback);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public <B, K, V, R> Future<?> getAndSet(final B bucket, final K key, final V value, final AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
public <B, K, V, R> Future<?> getAndSet(final B bucket, final K key, final V value,
|
|
|
|
|
final AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
final List<Future<?>> futures = new ArrayList<Future<?>>();
|
|
|
|
|
try {
|
|
|
|
|
getWithMetaData(bucket, key, null, new AsyncKeyValueStoreOperation<Object, Object>() {
|
|
|
|
|
@@ -242,11 +267,14 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck
|
|
|
|
|
return futures.size() > 0 ? futures.get(0) : null;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public <B, K, R> Future<?> getAndSetAsBytes(B bucket, K key, byte[] value, AsyncKeyValueStoreOperation<byte[], R> callback) {
|
|
|
|
|
public <B, K, R> Future<?> getAndSetAsBytes(B bucket, K key, byte[] value,
|
|
|
|
|
AsyncKeyValueStoreOperation<byte[], R> callback) {
|
|
|
|
|
return getAndSet(bucket, key, value, callback);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public <B, K, V, T, R> Future<?> getAndSetAsType(final B bucket, final K key, final V value, final Class<T> requiredType, final AsyncKeyValueStoreOperation<T, R> callback) {
|
|
|
|
|
public <B, K, V, T, R> Future<?> getAndSetAsType(final B bucket, final K key, final V value,
|
|
|
|
|
final Class<T> requiredType,
|
|
|
|
|
final AsyncKeyValueStoreOperation<T, R> callback) {
|
|
|
|
|
final List<Future<?>> futures = new ArrayList<Future<?>>();
|
|
|
|
|
getWithMetaData(bucket, key, requiredType, new AsyncKeyValueStoreOperation<T, R>() {
|
|
|
|
|
@SuppressWarnings({"unchecked"})
|
|
|
|
|
@@ -268,7 +296,8 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck
|
|
|
|
|
return futures.size() > 0 ? futures.get(0) : null;
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public <B, K, V, R> Future<?> setIfKeyNonExistent(final B bucket, final K key, final V value, final AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
public <B, K, V, R> Future<?> setIfKeyNonExistent(final B bucket, final K key, final V value,
|
|
|
|
|
final AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
return containsKey(bucket, key, new AsyncKeyValueStoreOperation<Boolean, Object>() {
|
|
|
|
|
public Object completed(KeyValueStoreMetaData meta, Boolean result) {
|
|
|
|
|
if (!result) {
|
|
|
|
|
@@ -284,7 +313,9 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck
|
|
|
|
|
});
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public <B, K, R> Future<?> setIfKeyNonExistentAsBytes(final B bucket, final K key, final byte[] value, final AsyncKeyValueStoreOperation<byte[], R> callback) {
|
|
|
|
|
public <B, K, R> Future<?> setIfKeyNonExistentAsBytes(final B bucket, final K key,
|
|
|
|
|
final byte[] value,
|
|
|
|
|
final AsyncKeyValueStoreOperation<byte[], R> callback) {
|
|
|
|
|
return containsKey(bucket, key, new AsyncKeyValueStoreOperation<Boolean, Object>() {
|
|
|
|
|
public Object completed(KeyValueStoreMetaData meta, Boolean result) {
|
|
|
|
|
if (!result) {
|
|
|
|
|
@@ -301,7 +332,8 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@SuppressWarnings({"unchecked"})
|
|
|
|
|
public <B, K, R> Future<?> containsKey(B bucket, K key, final AsyncKeyValueStoreOperation<Boolean, R> callback) {
|
|
|
|
|
public <B, K, R> Future<?> containsKey(B bucket, K key,
|
|
|
|
|
final AsyncKeyValueStoreOperation<Boolean, R> 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(),
|
|
|
|
|
@@ -318,24 +350,29 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
@SuppressWarnings({"unchecked"})
|
|
|
|
|
public <B, K, R> Future<?> delete(B bucket, K key, AsyncKeyValueStoreOperation<Boolean, R> callback) {
|
|
|
|
|
public <B, K, R> Future<?> delete(B bucket, K key,
|
|
|
|
|
AsyncKeyValueStoreOperation<Boolean, R> 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 <B, K, R> Future<?> setAsBytes(B bucket, K key, byte[] value, QosParameters qosParams, AsyncKeyValueStoreOperation<byte[], R> callback) {
|
|
|
|
|
public <B, K, R> Future<?> setAsBytes(B bucket, K key, byte[] value, QosParameters qosParams,
|
|
|
|
|
AsyncKeyValueStoreOperation<byte[], R> callback) {
|
|
|
|
|
return setWithMetaData(bucket, key, value, null, qosParams, callback);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
public <B, K, V, R> Future<?> setWithMetaData(B bucket, K key, V value, Map<String, String> metaData, AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
public <B, K, V, R> Future<?> setWithMetaData(B bucket, K key, V value,
|
|
|
|
|
Map<String, String> metaData,
|
|
|
|
|
AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
return setWithMetaData(bucket, key, value, metaData, null, callback);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/* ---------------- Map/Reduce ---------------- */
|
|
|
|
|
|
|
|
|
|
@SuppressWarnings({"unchecked"})
|
|
|
|
|
public <R> Future<?> execute(MapReduceJob job, AsyncKeyValueStoreOperation<List<?>, R> callback) {
|
|
|
|
|
public <R> Future<?> execute(MapReduceJob job,
|
|
|
|
|
AsyncKeyValueStoreOperation<List<?>, R> callback) {
|
|
|
|
|
HttpHeaders headers = defaultHeaders(null);
|
|
|
|
|
headers.setContentType(MediaType.APPLICATION_JSON);
|
|
|
|
|
HttpEntity<String> json = new HttpEntity<String>(job.toJson(), headers);
|
|
|
|
|
@@ -343,13 +380,15 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/* ---------------- Runnable helpers ---------------- */
|
|
|
|
|
|
|
|
|
|
protected class AsyncPut<V, R> implements Callable {
|
|
|
|
|
|
|
|
|
|
private String bucket;
|
|
|
|
|
private HttpEntity<V> entity = null;
|
|
|
|
|
private AsyncKeyValueStoreOperation<V, R> callback = null;
|
|
|
|
|
|
|
|
|
|
public AsyncPut(String bucket, HttpEntity<V> entity, AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
public AsyncPut(String bucket, HttpEntity<V> entity,
|
|
|
|
|
AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
this.bucket = bucket;
|
|
|
|
|
this.entity = entity;
|
|
|
|
|
this.callback = callback;
|
|
|
|
|
@@ -388,7 +427,8 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck
|
|
|
|
|
private HttpEntity<V> entity = null;
|
|
|
|
|
private AsyncKeyValueStoreOperation<V, R> callback = null;
|
|
|
|
|
|
|
|
|
|
public AsyncPost(String bucket, String key, HttpEntity<V> entity, AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
public AsyncPost(String bucket, String key, HttpEntity<V> entity,
|
|
|
|
|
AsyncKeyValueStoreOperation<V, R> callback) {
|
|
|
|
|
this.bucket = bucket;
|
|
|
|
|
this.key = key;
|
|
|
|
|
this.entity = entity;
|
|
|
|
|
@@ -433,7 +473,8 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck
|
|
|
|
|
private HttpEntity<String> entity = null;
|
|
|
|
|
private AsyncKeyValueStoreOperation<List<?>, R> callback = null;
|
|
|
|
|
|
|
|
|
|
public AsyncMapReduce(HttpEntity<String> entity, AsyncKeyValueStoreOperation<List<?>, R> callback) {
|
|
|
|
|
public AsyncMapReduce(HttpEntity<String> entity,
|
|
|
|
|
AsyncKeyValueStoreOperation<List<?>, R> callback) {
|
|
|
|
|
this.entity = entity;
|
|
|
|
|
this.callback = callback;
|
|
|
|
|
}
|
|
|
|
|
@@ -471,7 +512,8 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck
|
|
|
|
|
private Class<T> requiredType;
|
|
|
|
|
private AsyncKeyValueStoreOperation<T, R> callback = null;
|
|
|
|
|
|
|
|
|
|
public AsyncGet(String bucket, String key, Class<T> requiredType, AsyncKeyValueStoreOperation<T, R> callback) {
|
|
|
|
|
public AsyncGet(String bucket, String key, Class<T> requiredType,
|
|
|
|
|
AsyncKeyValueStoreOperation<T, R> callback) {
|
|
|
|
|
this.bucket = bucket;
|
|
|
|
|
this.key = key;
|
|
|
|
|
this.requiredType = requiredType;
|
|
|
|
|
@@ -520,7 +562,8 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck
|
|
|
|
|
private String key;
|
|
|
|
|
private AsyncKeyValueStoreOperation<HttpHeaders, R> callback = null;
|
|
|
|
|
|
|
|
|
|
public AsyncHead(String bucket, String key, AsyncKeyValueStoreOperation<HttpHeaders, R> callback) {
|
|
|
|
|
public AsyncHead(String bucket, String key,
|
|
|
|
|
AsyncKeyValueStoreOperation<HttpHeaders, R> callback) {
|
|
|
|
|
this.bucket = bucket;
|
|
|
|
|
this.key = key;
|
|
|
|
|
this.callback = callback;
|
|
|
|
|
@@ -552,7 +595,8 @@ public class AsyncRiakTemplate extends AbstractRiakTemplate implements AsyncBuck
|
|
|
|
|
private String key;
|
|
|
|
|
private AsyncKeyValueStoreOperation<Boolean, R> callback = null;
|
|
|
|
|
|
|
|
|
|
public AsyncDelete(String bucket, String key, AsyncKeyValueStoreOperation<Boolean, R> callback) {
|
|
|
|
|
public AsyncDelete(String bucket, String key,
|
|
|
|
|
AsyncKeyValueStoreOperation<Boolean, R> callback) {
|
|
|
|
|
this.bucket = bucket;
|
|
|
|
|
this.key = key;
|
|
|
|
|
this.callback = callback;
|
|
|
|
|
|