public class

RiakConnector

extends Object
implements BeanNameAware
java.lang.Object
   ↳ org.mule.modules.riak.RiakConnector
Known Direct Subclasses
Known Indirect Subclasses

Class Overview

Riak Connector. Based on Basho's Java client.

{@sample.config INCLUDE_ERROR} {@sample.config INCLUDE_ERROR} {@sample.config INCLUDE_ERROR} {@sample.config INCLUDE_ERROR} {@sample.config INCLUDE_ERROR}

Summary

Constants
String RIAK_FLOW_VAR_FETCHED_OBJECT
String RIAK_FLOW_VAR_PREFIX
Fields
private static final Logger LOGGER
private Configuration activeConfiguration
public int clusterTotalMaximumConnections If a cluster client is used (ie.
private RiakHttpClientConfigurationAdapter httpClientConfiguration The HTTP client configuration to use to connect to Riak.
private List<RiakHttpClientConfigurationAdapter> httpClientConfigurations The HTTP cluster client configurations to use to connect to Riak.
private boolean lazyLoadBucketProperties Allows to defer fetching bucket properties from Riak until they are required by one of the Bucket methods that accesses them.
private String name
private RiakProtobufClientConfigurationAdapter protobufClientConfiguration The Protocol Buffer client configuration to use to connect to Riak.
private List<RiakProtobufClientConfigurationAdapter> protobufClientConfigurations The Protocol Buffer cluster client configurations to use to connect to Riak.
private Retrier retrier A Retrier to use to perform actions.
private int retryCount The number of retries to attempt if no retrier has been configured and the DefaultRetrier is used.
private IRiakClient riakClient
Public Constructors
RiakConnector()
Public Methods
void connect()
Open the connection to Riak.
Bucket createBucket(String bucketName, boolean allowSiblings, String backend, boolean enableSearch, Boolean lastWriteWins, Boolean notFoundOK, int nVal, Integer smallVClock, Integer bigVClock, Long oldVClock, Long youngVClock, List<FunctionConfiguration> preCommitHooks, List<FunctionConfiguration> postCommitHooks, FunctionConfiguration chashKeyFunction, FunctionConfiguration linkWalkFunction, QuorumConfiguration quorumConfiguration)
Create a new bucket using the provided name.
void delete(String bucketName, String key, Object vClock, Boolean fetchBeforeDelete, QuorumConfiguration quorumConfiguration, MuleEvent muleEvent)
Delete the object at the specified key of the specified bucket.
void disconnect()
Close the connection to Riak.
byte[] fetch(String bucketName, String key, Object ifModifiedVClock, Date modifiedSince, Boolean notFoundOK, Boolean returnDeletedVClock, QuorumConfiguration quorumConfiguration, ConflictResolver<IRiakObject> resolver, MuleEvent muleEvent)
Fetches data from a Bucket and stores it in the current message payload.
List<String> fetchBinIndex(String bucketName, String indexName, String from, String to, String value)
Retrieve all the keys matching the bin index range or value.
Bucket fetchBucket(String bucketName)
Fetches a bucket by name.
List<String> fetchIntIndex(String bucketName, String indexName, Integer from, Integer to, Integer value)
Retrieve all the keys matching the integer index range or value.
Iterable<String> fetchKeys(String bucketName)
Retrieve all the keys in the specified bucket.
Configuration getActiveConfiguration()
int getClusterTotalMaximumConnections()
RiakHttpClientConfigurationAdapter getHttpClientConfiguration()
List<RiakHttpClientConfigurationAdapter> getHttpClientConfigurations()
String getName()
RiakProtobufClientConfigurationAdapter getProtobufClientConfiguration()
List<RiakProtobufClientConfigurationAdapter> getProtobufClientConfigurations()
Retrier getRetrier()
int getRetryCount()
IRiakClient getRiakClient()
Iterable<NodeStats> getStatistics()
Perform the Riak /stats operation on the node(s) this client is connected to.
boolean isLazyLoadBucketProperties()
Set<String> listBucketNames()
Retrieve the bucket names.
String mapReduceBinIndex(String bucketName, String indexName, String from, String to, String value, List<MapReducePhaseConfiguration> phases, Long timeout)
Performs a map-reduce operation that uses a Bin index query as input.
String mapReduceBucket(String bucketName, List<KeyFilter> keyFilters, List<MapReducePhaseConfiguration> phases, Long timeout)
Performs a map-reduce over a bucket.
String mapReduceBucketKey(List<BucketKeyInputConfiguration> inputs, List<MapReducePhaseConfiguration> phases, Long timeout)
Performs a map-reduce over a set of bucket/key/keydata inputs.
String mapReduceIntIndex(String bucketName, String indexName, Long from, Long to, Long value, List<MapReducePhaseConfiguration> phases, Long timeout)
Performs a map-reduce operation that uses an Int index query as input.
String mapReduceSearch(String bucketName, String query, List<MapReducePhaseConfiguration> phases, Long timeout)
Performs a map-reduce operation that uses a Riak Search query as input.
void setBeanName(String name)
void setClusterTotalMaximumConnections(int clusterTotalMaximumConnections)
void setHttpClientConfiguration(RiakHttpClientConfigurationAdapter httpClientConfiguration)
void setHttpClientConfigurations(List<RiakHttpClientConfigurationAdapter> httpClientConfigurations)
void setLazyLoadBucketProperties(boolean lazyLoadBucketProperties)
void setProtobufClientConfiguration(RiakProtobufClientConfigurationAdapter protobufClientConfiguration)
void setProtobufClientConfigurations(List<RiakProtobufClientConfigurationAdapter> protobufClientConfigurations)
void setRetrier(Retrier retrier)
void setRetryCount(int retryCount)
byte[] store(String bucketName, String key, String contentType, boolean ifNoneMatch, boolean ifNotModified, boolean returnBody, boolean withoutFetch, Boolean notFoundOK, Boolean returnDeletedVClock, QuorumConfiguration quorumConfiguration, ConflictResolver<IRiakObject> resolver, Mutation<IRiakObject> mutation, Object value, MuleEvent muleEvent)
Stores the current message payload into a Bucket.
Bucket updateBucket(String bucketName, boolean allowSiblings, String backend, boolean enableSearch, Boolean lastWriteWins, Boolean notFoundOK, int nVal, Integer smallVClock, Integer bigVClock, Long oldVClock, Long youngVClock, List<FunctionConfiguration> preCommitHooks, List<FunctionConfiguration> postCommitHooks, FunctionConfiguration chashKeyFunction, FunctionConfiguration linkWalkFunction, QuorumConfiguration quorumConfiguration)
Update a bucket using the provided name.
WalkResult walk(String bucketName, String key, List<LinkWalkStepConfiguration> linkWalks)
Walks the provided links starting at the object defined by the specified key and bucket name.
[Expand]
Inherited Methods
From class java.lang.Object
From interface org.springframework.beans.factory.BeanNameAware

Constants

public static final String RIAK_FLOW_VAR_FETCHED_OBJECT

Constant Value: "riak.fetchedObject"

public static final String RIAK_FLOW_VAR_PREFIX

Constant Value: "riak"

Fields

private static final Logger LOGGER

private Configuration activeConfiguration

public int clusterTotalMaximumConnections

If a cluster client is used (ie. if multiple client configurations are provided), defines the total of maximum connections. By default, no limit is set.

private RiakHttpClientConfigurationAdapter httpClientConfiguration

The HTTP client configuration to use to connect to Riak.

private List<RiakHttpClientConfigurationAdapter> httpClientConfigurations

The HTTP cluster client configurations to use to connect to Riak.

private boolean lazyLoadBucketProperties

Allows to defer fetching bucket properties from Riak until they are required by one of the Bucket methods that accesses them.

private String name

private RiakProtobufClientConfigurationAdapter protobufClientConfiguration

The Protocol Buffer client configuration to use to connect to Riak.

private List<RiakProtobufClientConfigurationAdapter> protobufClientConfigurations

The Protocol Buffer cluster client configurations to use to connect to Riak.

private Retrier retrier

A Retrier to use to perform actions. If none is configured, the DefaultRetrier will be used, configured for with retryCount retries.

private int retryCount

The number of retries to attempt if no retrier has been configured and the DefaultRetrier is used.

private IRiakClient riakClient

Public Constructors

public RiakConnector ()

Public Methods

public void connect ()

Open the connection to Riak.

Throws
RiakException if anything goes haywire when trying to reach Riak.
ConfigurationException if the configuration is invalid.
UnsupportedEncodingException

public Bucket createBucket (String bucketName, boolean allowSiblings, String backend, boolean enableSearch, Boolean lastWriteWins, Boolean notFoundOK, int nVal, Integer smallVClock, Integer bigVClock, Long oldVClock, Long youngVClock, List<FunctionConfiguration> preCommitHooks, List<FunctionConfiguration> postCommitHooks, FunctionConfiguration chashKeyFunction, FunctionConfiguration linkWalkFunction, QuorumConfiguration quorumConfiguration)

Create a new bucket using the provided name.

Parameters
bucketName The name of the new bucket.
allowSiblings Should the bucket have allow_mult set to true?
backend Which backend this bucket uses. Not supported by the Protobuf API.
enableSearch To enable or disable search and related commit hooks (support for both pre-1.0 and 1.0 search). Not supported by the Protobuf API.
lastWriteWins Set this bucket last_write_wins. Not supported by the Protobuf API.
notFoundOK Default notfound_ok value. Not supported by the Protobuf API.
nVal The n_val for this bucket.
smallVClock The small_vclock prune size. Not supported by the Protobuf API.
bigVClock The big_vclock prune size. Not supported by the Protobuf API.
oldVClock The old_vclock prune age. Not supported by the Protobuf API.
youngVClock The young_vclock prune age. Not supported by the Protobuf API.
preCommitHooks Erlang or JavaScript function names to be used as pre-commit hooks. Not supported by the Protobuf API.
postCommitHooks Erlang function names to be used as post-commit hooks. Not supported by the Protobuf API.
chashKeyFunction Erlang function names to be used as the chash_key_fun. Not supported by the Protobuf API.
linkWalkFunction Erlang function names to be used as the link_walk_fun. Not supported by the Protobuf API.
quorumConfiguration The specific quorum configuration to use for the operation.
Returns
  • the newly created Bucket.
Throws
RiakRetryFailedException thrown if the operation can't succeed.

public void delete (String bucketName, String key, Object vClock, Boolean fetchBeforeDelete, QuorumConfiguration quorumConfiguration, MuleEvent muleEvent)

Delete the object at the specified key of the specified bucket.

Parameters
bucketName The name of the bucket to fetch from.
key The key at which the data must be stored.
vClock Vector clock to delete. Acceptable types are VClock, byte[] or String.
fetchBeforeDelete If you want to provide a vclock to delete, but don't have one, setting this true will have the operation first perform a fetch (using the supplied r/pr parameters).
quorumConfiguration The specific quorum configuration to use for the operation.
muleEvent The current MuleEvent.
Throws
RiakException thrown if the operation can't succeed.
UnsupportedEncodingException thrown if a vector clock String can't be converted to byte[] using the MuleEvent's encoding.

public void disconnect ()

Close the connection to Riak.

public byte[] fetch (String bucketName, String key, Object ifModifiedVClock, Date modifiedSince, Boolean notFoundOK, Boolean returnDeletedVClock, QuorumConfiguration quorumConfiguration, ConflictResolver<IRiakObject> resolver, MuleEvent muleEvent)

Fetches data from a Bucket and stores it in the current message payload.

Parameters
bucketName The name of the bucket to fetch from.
key The key at which the data must be stored.
ifModifiedVClock Fetch only if the stored vector clock is different from the provided vector clock. Acceptable types are VClock , byte[] or String.
modifiedSince Fetch only if the stored value has been modified after the provided date.
notFoundOK If a notfound response counts towards satisfying the r value.
returnDeletedVClock If an object has been deleted, return the tombstone vclock.
quorumConfiguration The specific quorum configuration to use for the operation.
resolver The ConflictResolver to use on any sibling results returned from the fetch (and store if returnBody is true).
muleEvent The current MuleEvent.
Returns
  • the byte[] representing the stored value if returnBody is true, null if returnBody is false.
Throws
RiakRetryFailedException thrown if the operation can't succeed.
UnsupportedEncodingException thrown if a vector clock String can't be converted to byte[] using the MuleEvent's encoding.

public List<String> fetchBinIndex (String bucketName, String indexName, String from, String to, String value)

Retrieve all the keys matching the bin index range or value.

Parameters
bucketName The name of the bucket to fetch from.
indexName The name of the index.
from The start value of the index range.
to The end value of the index range.
value The fixed index value.
Returns
  • the list of keys that match the index range or value.
Throws
RiakException thrown if the operation can't succeed.

public Bucket fetchBucket (String bucketName)

Fetches a bucket by name.

Parameters
bucketName The name of the bucket.
Returns
  • the fetched Bucket.
Throws
RiakRetryFailedException thrown if the operation can't succeed.

public List<String> fetchIntIndex (String bucketName, String indexName, Integer from, Integer to, Integer value)

Retrieve all the keys matching the integer index range or value.

Parameters
bucketName The name of the bucket to fetch from.
indexName The name of the index.
from The start value of the index range.
to The end value of the index range.
value The fixed index value.
Returns
  • the list of keys that match the index range or value.
Throws
RiakException thrown if the operation can't succeed.

public Iterable<String> fetchKeys (String bucketName)

Retrieve all the keys in the specified bucket.

Parameters
bucketName The name of the bucket to fetch from.
Returns
  • an Iterable of String keys.
Throws
RiakException thrown if the operation can't succeed.

public Configuration getActiveConfiguration ()

Throws
ConfigurationException

public int getClusterTotalMaximumConnections ()

public RiakHttpClientConfigurationAdapter getHttpClientConfiguration ()

public List<RiakHttpClientConfigurationAdapter> getHttpClientConfigurations ()

public String getName ()

public RiakProtobufClientConfigurationAdapter getProtobufClientConfiguration ()

public List<RiakProtobufClientConfigurationAdapter> getProtobufClientConfigurations ()

public Retrier getRetrier ()

public int getRetryCount ()

public IRiakClient getRiakClient ()

public Iterable<NodeStats> getStatistics ()

Perform the Riak /stats operation on the node(s) this client is connected to.

This is not supported by the Protobuf API.

Returns
  • an Iterable object that contains one or more NodeStats
Throws
RiakException If Riak does not respond or if the protobuf API is being used

public boolean isLazyLoadBucketProperties ()

public Set<String> listBucketNames ()

Retrieve the bucket names.

Returns
  • a set of Bucket names.
Throws
RiakException thrown if the fetch operation can't succeed.

public String mapReduceBinIndex (String bucketName, String indexName, String from, String to, String value, List<MapReducePhaseConfiguration> phases, Long timeout)

Performs a map-reduce operation that uses a Bin index query as input.

Parameters
bucketName The name of the bucket to execute the map-reduce on.
indexName The name of the index to query.
from The start value of the index range.
to The end value of the index range.
value The fixed index value.
phases A List of map-reduce phases defined with MapReducePhaseConfiguration instances.
timeout The maximum duration in milliseconds for the map-reduce operation.
Returns
  • the raw JSON string of the result.
Throws
RiakException thrown if the operation can't succeed.

public String mapReduceBucket (String bucketName, List<KeyFilter> keyFilters, List<MapReducePhaseConfiguration> phases, Long timeout)

Performs a map-reduce over a bucket.

Parameters
bucketName The name of the bucket to execute the map-reduce on.
keyFilters A List of KeyFilter to apply on the inputs of the map-reduce operation.
phases A List of map-reduce phases defined with MapReducePhaseConfiguration instances.
timeout The maximum duration in milliseconds for the map-reduce operation.
Returns
  • the raw JSON string of the result.
Throws
RiakException thrown if the operation can't succeed.

public String mapReduceBucketKey (List<BucketKeyInputConfiguration> inputs, List<MapReducePhaseConfiguration> phases, Long timeout)

Performs a map-reduce over a set of bucket/key/keydata inputs.

Parameters
inputs A List of BucketKeyInputConfiguration used as input for the map-reduce operation.
phases A List of map-reduce phases defined with MapReducePhaseConfiguration instances.
timeout The maximum duration in milliseconds for the map-reduce operation.
Returns
  • the raw JSON string of the result.
Throws
RiakException thrown if the operation can't succeed.

public String mapReduceIntIndex (String bucketName, String indexName, Long from, Long to, Long value, List<MapReducePhaseConfiguration> phases, Long timeout)

Performs a map-reduce operation that uses an Int index query as input.

Parameters
bucketName The name of the bucket to execute the map-reduce on.
indexName The name of the index to query.
from The start value of the index range.
to The end value of the index range.
value The fixed index value.
phases A List of map-reduce phases defined with MapReducePhaseConfiguration instances.
timeout The maximum duration in milliseconds for the map-reduce operation.
Returns
  • the raw JSON string of the result.
Throws
RiakException thrown if the operation can't succeed.

public String mapReduceSearch (String bucketName, String query, List<MapReducePhaseConfiguration> phases, Long timeout)

Performs a map-reduce operation that uses a Riak Search query as input.

Parameters
bucketName The name of the bucket to execute the map-reduce on.
query The query to run to provide data to the map-reduce operation.
phases A List of map-reduce phases defined with MapReducePhaseConfiguration instances.
timeout The maximum duration in milliseconds for the map-reduce operation.
Returns
  • the raw JSON string of the result.
Throws
RiakException thrown if the operation can't succeed.

public void setBeanName (String name)

Parameters
name

public void setClusterTotalMaximumConnections (int clusterTotalMaximumConnections)

Parameters
clusterTotalMaximumConnections

public void setHttpClientConfiguration (RiakHttpClientConfigurationAdapter httpClientConfiguration)

Parameters
httpClientConfiguration

public void setHttpClientConfigurations (List<RiakHttpClientConfigurationAdapter> httpClientConfigurations)

Parameters
httpClientConfigurations

public void setLazyLoadBucketProperties (boolean lazyLoadBucketProperties)

Parameters
lazyLoadBucketProperties

public void setProtobufClientConfiguration (RiakProtobufClientConfigurationAdapter protobufClientConfiguration)

Parameters
protobufClientConfiguration

public void setProtobufClientConfigurations (List<RiakProtobufClientConfigurationAdapter> protobufClientConfigurations)

Parameters
protobufClientConfigurations

public void setRetrier (Retrier retrier)

Parameters
retrier

public void setRetryCount (int retryCount)

Parameters
retryCount

public byte[] store (String bucketName, String key, String contentType, boolean ifNoneMatch, boolean ifNotModified, boolean returnBody, boolean withoutFetch, Boolean notFoundOK, Boolean returnDeletedVClock, QuorumConfiguration quorumConfiguration, ConflictResolver<IRiakObject> resolver, Mutation<IRiakObject> mutation, Object value, MuleEvent muleEvent)

Stores the current message payload into a Bucket.

Parameters
bucketName The name of the bucket to store into.
key The key at which the data must be stored, which is optional if the current message payload is a IRiakObject.
contentType The content type of the data being stored. Defaults to the mime type of the current message payload.
ifNoneMatch True if you want a conditional store, false otherwise, defaults to false. NOTE: This has different meanings depending on the underlying transport.
ifNotModified True if you want a conditional store, false otherwise, defaults to false. NOTE: This has different meanings depending on the underlying transport.
returnBody Should the store operation return a response body?
withoutFetch Eliminates fetching the existing value before storing the current one.
notFoundOK If notfound_ok counts towards r count (for the pre-store fetch).
returnDeletedVClock If the object has just been deleted, there maybe a tombstone value vclock, set to true to have this returned in the pre-store fetch.
quorumConfiguration The specific quorum configuration to use for the operation.
resolver The ConflictResolver to use on any sibling results returned from the fetch (and store if returnBody is true).
mutation If provided, the current message payload is disregard and the Mutation is applied to generate the data to store.
value The current message payload to store, auto-transformed into bytes[] if needed, or stored as-is if it's a IRiakObject.
muleEvent The current MuleEvent.
Returns
  • the byte[] representing the stored value if returnBody is true, null if returnBody is false.
Throws
RiakRetryFailedException thrown if the operation can't succeed.
MuleException thrown if the operation can't succeed.

public Bucket updateBucket (String bucketName, boolean allowSiblings, String backend, boolean enableSearch, Boolean lastWriteWins, Boolean notFoundOK, int nVal, Integer smallVClock, Integer bigVClock, Long oldVClock, Long youngVClock, List<FunctionConfiguration> preCommitHooks, List<FunctionConfiguration> postCommitHooks, FunctionConfiguration chashKeyFunction, FunctionConfiguration linkWalkFunction, QuorumConfiguration quorumConfiguration)

Update a bucket using the provided name.

Parameters
bucketName The name of the new bucket.
allowSiblings Should the bucket have allow_mult set to true?
backend Which backend this bucket uses. Not supported by the Protobuf API.
enableSearch To enable or disable search and related commit hooks (support for both pre-1.0 and 1.0 search). Not supported by the Protobuf API.
lastWriteWins Set this bucket last_write_wins. Not supported by the Protobuf API.
notFoundOK Default notfound_ok value. Not supported by the Protobuf API.
nVal The n_val for this bucket.
smallVClock The small_vclock prune size. Not supported by the Protobuf API.
bigVClock The big_vclock prune size. Not supported by the Protobuf API.
oldVClock The old_vclock prune age. Not supported by the Protobuf API.
youngVClock The young_vclock prune age. Not supported by the Protobuf API.
preCommitHooks Erlang or JavaScript function names to be used as pre-commit hooks. Not supported by the Protobuf API.
postCommitHooks Erlang function names to be used as post-commit hooks. Not supported by the Protobuf API.
chashKeyFunction Erlang function names to be used as the chash_key_fun. Not supported by the Protobuf API.
linkWalkFunction Erlang function names to be used as the link_walk_fun. Not supported by the Protobuf API.
quorumConfiguration The specific quorum configuration to use for the operation.
Returns
  • the newly created Bucket.
Throws
RiakRetryFailedException thrown if the operation can't succeed.

public WalkResult walk (String bucketName, String key, List<LinkWalkStepConfiguration> linkWalks)

Walks the provided links starting at the object defined by the specified key and bucket name.

Parameters
bucketName The name of the bucket to fetch from.
key The key at which the data must be stored.
linkWalks The links that must be walked.
Returns
  • the WalkResult.
Throws
RiakException thrown if the operation can't succeed.