diff --git a/src/main/java/org/springframework/data/cassandra/core/CassandraOperations.java b/src/main/java/org/springframework/data/cassandra/core/CassandraOperations.java index 6ba3a4bcf..d9035f1cd 100644 --- a/src/main/java/org/springframework/data/cassandra/core/CassandraOperations.java +++ b/src/main/java/org/springframework/data/cassandra/core/CassandraOperations.java @@ -16,6 +16,7 @@ package org.springframework.data.cassandra.core; import java.util.List; +import java.util.Map; import org.springframework.data.cassandra.convert.CassandraConverter; @@ -25,6 +26,10 @@ import com.datastax.driver.core.ResultSetFuture; /** * @author Alex Shvid */ +/** + * @author David Webb + * + */ public interface CassandraOperations { /** @@ -92,6 +97,22 @@ public interface CassandraOperations { */ T insert(T entity, String tableName); + /** + * @param entity + * @param tableName + * @param options + * @return + */ + T insert(T entity, String tableName, QueryOptions options); + + /** + * @param entity + * @param tableName + * @param optionsByName + * @return + */ + T insert(T entity, String tableName, Map optionsByName); + /** * Insert the given list of objects to the table by annotation table name. * @@ -109,6 +130,22 @@ public interface CassandraOperations { */ List insert(List entities, String tableName); + /** + * @param entities + * @param tableName + * @param options + * @return + */ + List insert(List entities, String tableName, QueryOptions options); + + /** + * @param entities + * @param tableName + * @param optionsByName + * @return + */ + List insert(List entities, String tableName, Map optionsByName); + /** * Insert the given object to the table by id. * @@ -116,6 +153,29 @@ public interface CassandraOperations { */ T insertAsynchronously(T entity); + /** + * Insert the given object to the table by id. + * + * @param object + */ + T insertAsynchronously(T entity, String tableName); + + /** + * @param entity + * @param tableName + * @param options + * @return + */ + T insertAsynchronously(T entity, String tableName, QueryOptions options); + + /** + * @param entity + * @param tableName + * @param optionsByName + * @return + */ + T insertAsynchronously(T entity, String tableName, Map optionsByName); + /** * Insert the given object to the table by id. * @@ -128,14 +188,23 @@ public interface CassandraOperations { * * @param object */ - T insertAsynchronously(T entity, String tableName); + List insertAsynchronously(List entities, String tableName); /** - * Insert the given object to the table by id. - * - * @param object + * @param entities + * @param tableName + * @param options + * @return */ - List insertAsynchronously(List entities, String tableName); + List insertAsynchronously(List entities, String tableName, QueryOptions options); + + /** + * @param entities + * @param tableName + * @param optionsByName + * @return + */ + List insertAsynchronously(List entities, String tableName, Map optionsByName); /** * Insert the given object to the table by id. @@ -144,6 +213,29 @@ public interface CassandraOperations { */ T update(T entity); + /** + * Insert the given object to the table by id. + * + * @param object + */ + T update(T entity, String tableName); + + /** + * @param entity + * @param tableName + * @param options + * @return + */ + T update(T entity, String tableName, QueryOptions options); + + /** + * @param entity + * @param tableName + * @param optionsByName + * @return + */ + T update(T entity, String tableName, Map optionsByName); + /** * Insert the given object to the table by id. * @@ -156,14 +248,23 @@ public interface CassandraOperations { * * @param object */ - T update(T entity, String tableName); + List update(List entities, String tableName); /** - * Insert the given object to the table by id. - * - * @param object + * @param entities + * @param tableName + * @param options + * @return */ - List update(List entities, String tableName); + List update(List entities, String tableName, QueryOptions options); + + /** + * @param entities + * @param tableName + * @param optionsByName + * @return + */ + List update(List entities, String tableName, Map optionsByName); /** * Insert the given object to the table by id. @@ -177,14 +278,30 @@ public interface CassandraOperations { * * @param object */ - List updateAsynchronously(List entities); + T updateAsynchronously(T entity, String tableName); + + /** + * @param entity + * @param tableName + * @param options + * @return + */ + T updateAsynchronously(T entity, String tableName, QueryOptions options); + + /** + * @param entity + * @param tableName + * @param optionsByName + * @return + */ + T updateAsynchronously(T entity, String tableName, Map optionsByName); /** * Insert the given object to the table by id. * * @param object */ - T updateAsynchronously(T entity, String tableName); + List updateAsynchronously(List entities); /** * Insert the given object to the table by id. @@ -193,6 +310,22 @@ public interface CassandraOperations { */ List updateAsynchronously(List entities, String tableName); + /** + * @param entities + * @param tableName + * @param options + * @return + */ + List updateAsynchronously(List entities, String tableName, QueryOptions options); + + /** + * @param entities + * @param tableName + * @param optionsByName + * @return + */ + List updateAsynchronously(List entities, String tableName, Map optionsByName); + /** * Remove the given object from the table by id. * @@ -200,6 +333,28 @@ public interface CassandraOperations { */ void delete(T entity); + /** + * Removes the given object from the given table. + * + * @param object + * @param table must not be {@literal null} or empty. + */ + void delete(T entity, String tableName); + + /** + * @param entity + * @param tableName + * @param options + */ + void delete(T entity, String tableName, QueryOptions options); + + /** + * @param entity + * @param tableName + * @param optionsByName + */ + void delete(T entity, String tableName, Map optionsByName); + /** * Remove the given object from the table by id. * @@ -213,15 +368,21 @@ public interface CassandraOperations { * @param object * @param table must not be {@literal null} or empty. */ - void delete(T entity, String tableName); + void delete(List entities, String tableName); /** - * Removes the given object from the given table. - * - * @param object - * @param table must not be {@literal null} or empty. + * @param entities + * @param tableName + * @param options */ - void delete(List entities, String tableName); + void delete(List entities, String tableName, QueryOptions options); + + /** + * @param entities + * @param tableName + * @param optionsByName + */ + void delete(List entities, String tableName, Map optionsByName); /** * Remove the given object from the table by id. @@ -230,6 +391,28 @@ public interface CassandraOperations { */ void deleteAsychronously(T entity); + /** + * @param entity + * @param tableName + * @param options + */ + void deleteAsychronously(T entity, String tableName, QueryOptions options); + + /** + * @param entity + * @param tableName + * @param optionsByName + */ + void deleteAsychronously(T entity, String tableName, Map optionsByName); + + /** + * Removes the given object from the given table. + * + * @param object + * @param table must not be {@literal null} or empty. + */ + void deleteAsychronously(T entity, String tableName); + /** * Remove the given object from the table by id. * @@ -243,15 +426,21 @@ public interface CassandraOperations { * @param object * @param table must not be {@literal null} or empty. */ - void deleteAsychronously(T entity, String tableName); + void deleteAsychronously(List entities, String tableName); /** - * Removes the given object from the given table. - * - * @param object - * @param table must not be {@literal null} or empty. + * @param entities + * @param tableName + * @param options */ - void deleteAsychronously(List entities, String tableName); + void deleteAsychronously(List entities, String tableName, QueryOptions options); + + /** + * @param entities + * @param tableName + * @param optionsByName + */ + void deleteAsychronously(List entities, String tableName, Map optionsByName); /** * Returns the underlying {@link CassandraConverter}. diff --git a/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java b/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java index 2a11ed277..92314f1cd 100644 --- a/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java +++ b/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java @@ -18,9 +18,11 @@ package org.springframework.data.cassandra.core; import java.util.ArrayList; import java.util.Collection; import java.util.Collections; +import java.util.HashMap; import java.util.HashSet; import java.util.Iterator; import java.util.List; +import java.util.Map; import java.util.Set; import org.slf4j.Logger; @@ -57,8 +59,34 @@ import com.datastax.driver.core.querybuilder.Batch; */ public class CassandraTemplate implements CassandraOperations { + /** + * Simple {@link RowCallback} that will transform {@link Row} into the given target type using the given + * {@link EntityReader}. + * + * @author Alex Shvid + */ + private static class ReadRowCallback implements RowCallback { + + private final EntityReader reader; + private final Class type; + + public ReadRowCallback(EntityReader reader, Class type) { + Assert.notNull(reader); + Assert.notNull(type); + this.reader = reader; + this.type = type; + } + + @Override + public T doWith(Row object) { + T source = reader.read(type, object); + return source; + } + } + private static Logger log = LoggerFactory.getLogger(CassandraTemplate.class); public static final Collection ITERABLE_CLASSES; + static { Set iterableClasses = new HashSet(); @@ -69,12 +97,12 @@ public class CassandraTemplate implements CassandraOperations { ITERABLE_CLASSES = Collections.unmodifiableCollection(iterableClasses); } - private final Keyspace keyspace; private final Session session; private final CassandraConverter cassandraConverter; private final MappingContext, CassandraPersistentProperty> mappingContext; private final PersistenceExceptionTranslator exceptionTranslator = new CassandraExceptionTranslator(); + private ClassLoader beanClassLoader; /** @@ -89,11 +117,154 @@ public class CassandraTemplate implements CassandraOperations { this.mappingContext = this.cassandraConverter.getMappingContext(); } - /** - * @param classLoader + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#delete(java.util.List) */ - public void setBeanClassLoader(ClassLoader classLoader) { - this.beanClassLoader = classLoader; + @Override + public void delete(List entities) { + String tableName = getTableName(entities.get(0).getClass()); + Assert.notNull(tableName); + delete(entities, tableName); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#delete(java.util.List, java.lang.String) + */ + @Override + public void delete(List entities, String tableName) { + delete(entities, tableName, new HashMap()); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#delete(java.util.List, java.lang.String, java.util.Map) + */ + @Override + public void delete(List entities, String tableName, Map optionsByName) { + Assert.notNull(entities); + Assert.notEmpty(entities); + Assert.notNull(tableName); + Assert.notNull(optionsByName); + doBatchDelete(tableName, entities, optionsByName, false); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#delete(java.util.List, java.lang.String, org.springframework.data.cassandra.core.QueryOptions) + */ + @Override + public void delete(List entities, String tableName, QueryOptions options) { + delete(entities, tableName, options.toMap()); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#delete(java.lang.Object) + */ + @Override + public void delete(T entity) { + String tableName = getTableName(entity.getClass()); + Assert.notNull(tableName); + delete(entity, tableName); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#delete(java.lang.Object, java.lang.String) + */ + @Override + public void delete(T entity, String tableName) { + delete(entity, tableName, new HashMap()); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#delete(java.lang.Object, java.lang.String, java.util.Map) + */ + @Override + public void delete(T entity, String tableName, Map optionsByName) { + Assert.notNull(entity); + Assert.notNull(tableName); + Assert.notNull(optionsByName); + doDelete(tableName, entity, optionsByName, false); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#delete(java.lang.Object, java.lang.String, org.springframework.data.cassandra.core.QueryOptions) + */ + @Override + public void delete(T entity, String tableName, QueryOptions options) { + delete(entity, tableName, options.toMap()); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#deleteAsychronously(java.util.List) + */ + @Override + public void deleteAsychronously(List entities) { + String tableName = getTableName(entities.get(0).getClass()); + Assert.notNull(tableName); + deleteAsychronously(entities, tableName); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#deleteAsychronously(java.util.List, java.lang.String) + */ + @Override + public void deleteAsychronously(List entities, String tableName) { + deleteAsychronously(entities, tableName, new HashMap()); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#deleteAsychronously(java.util.List, java.lang.String, java.util.Map) + */ + @Override + public void deleteAsychronously(List entities, String tableName, Map optionsByName) { + Assert.notNull(entities); + Assert.notEmpty(entities); + Assert.notNull(tableName); + Assert.notNull(optionsByName); + doBatchDelete(tableName, entities, optionsByName, true); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#deleteAsychronously(java.util.List, java.lang.String, org.springframework.data.cassandra.core.QueryOptions) + */ + @Override + public void deleteAsychronously(List entities, String tableName, QueryOptions options) { + deleteAsychronously(entities, tableName, options.toMap()); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#deleteAsychronously(java.lang.Object) + */ + @Override + public void deleteAsychronously(T entity) { + String tableName = getTableName(entity.getClass()); + Assert.notNull(tableName); + deleteAsychronously(entity, tableName); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#deleteAsychronously(java.lang.Object, java.lang.String) + */ + @Override + public void deleteAsychronously(T entity, String tableName) { + deleteAsychronously(entity, tableName, new HashMap()); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#deleteAsychronously(java.lang.Object, java.lang.String, java.util.Map) + */ + @Override + public void deleteAsychronously(T entity, String tableName, Map optionsByName) { + Assert.notNull(entity); + Assert.notNull(tableName); + Assert.notNull(optionsByName); + doDelete(tableName, entity, optionsByName, true); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#deleteAsychronously(java.lang.Object, java.lang.String, org.springframework.data.cassandra.core.QueryOptions) + */ + @Override + public void deleteAsychronously(T entity, String tableName, QueryOptions options) { + deleteAsychronously(entity, tableName, options.toMap()); } /* (non-Javadoc) @@ -146,39 +317,28 @@ public class CassandraTemplate implements CassandraOperations { } /** - * Determines the PersistentEntityType for a given Object - * - * @param o + * @param entityClass * @return */ - protected CassandraPersistentEntity getEntity(Object o) { + public String determineTableName(Class entityClass) { - CassandraPersistentEntity entity = null; - try { - String entityClassName = o.getClass().getName(); - Class entityClass = ClassUtils.forName(entityClassName, beanClassLoader); - entity = mappingContext.getPersistentEntity(entityClass); - } catch (ClassNotFoundException e) { - e.printStackTrace(); - } catch (LinkageError e) { - e.printStackTrace(); - } finally { + if (entityClass == null) { + throw new InvalidDataAccessApiUsageException( + "No class parameter provided, entity table name can't be determined!"); } - return entity; - - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#getTableName(java.lang.Class) - */ - public String getTableName(Class entityClass) { - return determineTableName(entityClass); + CassandraPersistentEntity entity = mappingContext.getPersistentEntity(entityClass); + if (entity == null) { + throw new InvalidDataAccessApiUsageException("No Persitent Entity information found for the class " + + entityClass.getName()); + } + return entity.getTable(); } /* (non-Javadoc) * @see org.springframework.data.cassandra.core.CassandraOperations#executeQuery(java.lang.String) */ + @Override public ResultSet executeQuery(final String query) { return execute(new SessionCallback() { @@ -213,9 +373,179 @@ public class CassandraTemplate implements CassandraOperations { } + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#getConverter() + */ + @Override + public CassandraConverter getConverter() { + return cassandraConverter; + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#getTableName(java.lang.Class) + */ + @Override + public String getTableName(Class entityClass) { + return determineTableName(entityClass); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#insert(java.util.List) + */ + @Override + public List insert(List entities) { + String tableName = getTableName(entities.get(0).getClass()); + Assert.notNull(tableName); + return insert(entities, tableName); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#insert(java.util.List, java.lang.String) + */ + @Override + public List insert(List entities, String tableName) { + return insert(entities, tableName, new HashMap()); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#insert(java.util.List, java.lang.String, java.util.Map) + */ + @Override + public List insert(List entities, String tableName, Map optionsByName) { + Assert.notNull(entities); + Assert.notEmpty(entities); + Assert.notNull(tableName); + Assert.notNull(optionsByName); + return doBatchInsert(tableName, entities, optionsByName, false); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#insert(java.util.List, java.lang.String, org.springframework.data.cassandra.core.QueryOptions) + */ + @Override + public List insert(List entities, String tableName, QueryOptions options) { + return insert(entities, tableName, options.toMap()); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#insert(java.lang.Object) + */ + @Override + public T insert(T entity) { + String tableName = determineTableName(entity); + Assert.notNull(tableName); + return insert(entity, tableName); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#insert(java.lang.Object, java.lang.String) + */ + @Override + public T insert(T entity, String tableName) { + return insert(entity, tableName, new HashMap()); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#insert(java.lang.Object, java.lang.String, java.util.Map) + */ + @Override + public T insert(T entity, String tableName, Map optionsByName) { + Assert.notNull(entity); + Assert.notNull(tableName); + ensureNotIterable(entity); + return doInsert(tableName, entity, optionsByName, false); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#insert(java.lang.Object, java.lang.String, org.springframework.data.cassandra.core.QueryOptions) + */ + @Override + public T insert(T entity, String tableName, QueryOptions options) { + return insert(entity, tableName, options.toMap()); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#insertAsynchronously(java.util.List) + */ + @Override + public List insertAsynchronously(List entities) { + String tableName = getTableName(entities.get(0).getClass()); + Assert.notNull(tableName); + return insertAsynchronously(entities, tableName); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#insertAsynchronously(java.util.List, java.lang.String) + */ + @Override + public List insertAsynchronously(List entities, String tableName) { + return insertAsynchronously(entities, tableName, new HashMap()); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#insertAsynchronously(java.util.List, java.lang.String, java.util.Map) + */ + @Override + public List insertAsynchronously(List entities, String tableName, Map optionsByName) { + Assert.notNull(entities); + Assert.notEmpty(entities); + Assert.notNull(tableName); + Assert.notNull(optionsByName); + return doBatchInsert(tableName, entities, optionsByName, true); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#insertAsynchronously(java.util.List, java.lang.String, org.springframework.data.cassandra.core.QueryOptions) + */ + @Override + public List insertAsynchronously(List entities, String tableName, QueryOptions options) { + return insertAsynchronously(entities, tableName, options.toMap()); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#insertAsynchronously(java.lang.Object) + */ + @Override + public T insertAsynchronously(T entity) { + String tableName = determineTableName(entity); + Assert.notNull(tableName); + return insertAsynchronously(entity, tableName); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#insertAsynchronously(java.lang.Object, java.lang.String) + */ + @Override + public T insertAsynchronously(T entity, String tableName) { + return insertAsynchronously(entity, tableName, new HashMap()); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#insertAsynchronously(java.lang.Object, java.lang.String, java.util.Map) + */ + @Override + public T insertAsynchronously(T entity, String tableName, Map optionsByName) { + Assert.notNull(entity); + Assert.notNull(tableName); + Assert.notNull(optionsByName); + + ensureNotIterable(entity); + + return doInsert(tableName, entity, optionsByName, true); + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#insertAsynchronously(java.lang.Object, java.lang.String, org.springframework.data.cassandra.core.QueryOptions) + */ + @Override + public T insertAsynchronously(T entity, String tableName, QueryOptions options) { + return insertAsynchronously(entity, tableName, options.toMap()); + } + /* (non-Javadoc) * @see org.springframework.data.cassandra.core.CassandraOperations#select(java.lang.String, java.lang.Class) */ + @Override public List select(String query, Class selectClass) { return selectInternal(query, new ReadRowCallback(cassandraConverter, selectClass)); } @@ -223,41 +553,400 @@ public class CassandraTemplate implements CassandraOperations { /* (non-Javadoc) * @see org.springframework.data.cassandra.core.CassandraOperations#selectOne(java.lang.String, java.lang.Class) */ + @Override public T selectOne(String query, Class selectClass) { return selectOneInternal(query, new ReadRowCallback(cassandraConverter, selectClass)); } - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#getConverter() + /** + * @param classLoader */ - public CassandraConverter getConverter() { - return cassandraConverter; + public void setBeanClassLoader(ClassLoader classLoader) { + this.beanClassLoader = classLoader; + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#update(java.util.List) + */ + @Override + public List update(List entities) { + // TODO Auto-generated method stub + return null; + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#update(java.util.List, java.lang.String) + */ + @Override + public List update(List entities, String tableName) { + // TODO Auto-generated method stub + return null; + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#update(java.util.List, java.lang.String, java.util.Map) + */ + @Override + public List update(List entities, String tableName, Map optionsByName) { + // TODO Auto-generated method stub + return null; + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#update(java.util.List, java.lang.String, org.springframework.data.cassandra.core.QueryOptions) + */ + @Override + public List update(List entities, String tableName, QueryOptions options) { + // TODO Auto-generated method stub + return null; + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#update(java.lang.Object) + */ + @Override + public T update(T entity) { + // TODO Auto-generated method stub + return null; + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#update(java.lang.Object, java.lang.String) + */ + @Override + public T update(T entity, String tableName) { + // TODO Auto-generated method stub + return null; + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#update(java.lang.Object, java.lang.String, java.util.Map) + */ + @Override + public T update(T entity, String tableName, Map optionsByName) { + // TODO Auto-generated method stub + return null; + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#update(java.lang.Object, java.lang.String, org.springframework.data.cassandra.core.QueryOptions) + */ + @Override + public T update(T entity, String tableName, QueryOptions options) { + // TODO Auto-generated method stub + return null; + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#updateAsynchronously(java.util.List) + */ + @Override + public List updateAsynchronously(List entities) { + // TODO Auto-generated method stub + return null; + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#updateAsynchronously(java.util.List, java.lang.String) + */ + @Override + public List updateAsynchronously(List entities, String tableName) { + // TODO Auto-generated method stub + return null; + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#updateAsynchronously(java.util.List, java.lang.String, java.util.Map) + */ + @Override + public List updateAsynchronously(List entities, String tableName, Map optionsByName) { + // TODO Auto-generated method stub + return null; + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#updateAsynchronously(java.util.List, java.lang.String, org.springframework.data.cassandra.core.QueryOptions) + */ + @Override + public List updateAsynchronously(List entities, String tableName, QueryOptions options) { + // TODO Auto-generated method stub + return null; + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#updateAsynchronously(java.lang.Object) + */ + @Override + public T updateAsynchronously(T entity) { + // TODO Auto-generated method stub + return null; + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#updateAsynchronously(java.lang.Object, java.lang.String) + */ + @Override + public T updateAsynchronously(T entity, String tableName) { + // TODO Auto-generated method stub + return null; + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#updateAsynchronously(java.lang.Object, java.lang.String, java.util.Map) + */ + @Override + public T updateAsynchronously(T entity, String tableName, Map optionsByName) { + // TODO Auto-generated method stub + return null; + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#updateAsynchronously(java.lang.Object, java.lang.String, org.springframework.data.cassandra.core.QueryOptions) + */ + @Override + public T updateAsynchronously(T entity, String tableName, QueryOptions options) { + // TODO Auto-generated method stub + return null; } /** - * Simple {@link RowCallback} that will transform {@link Row} into the given target type using the given - * {@link EntityReader}. - * - * @author Alex Shvid + * @param obj + * @return */ - private static class ReadRowCallback implements RowCallback { - - private final EntityReader reader; - private final Class type; - - public ReadRowCallback(EntityReader reader, Class type) { - Assert.notNull(reader); - Assert.notNull(type); - this.reader = reader; - this.type = type; + private String determineTableName(T obj) { + if (null != obj) { + return determineTableName(obj.getClass()); } - public T doWith(Row object) { - T source = reader.read(type, object); - return source; + return null; + } + + private RuntimeException potentiallyConvertRuntimeException(RuntimeException ex) { + RuntimeException resolved = this.exceptionTranslator.translateExceptionIfPossible(ex); + return resolved == null ? ex : resolved; + } + + /** + * Perform the deletion on a list of objects + * + * @param tableName + * @param objectToRemove + */ + protected void doBatchDelete(final String tableName, final List entities, Map optionsByName, + final boolean deleteAsynchronously) { + + Assert.notEmpty(entities); + + CassandraPersistentEntity CPEntity = getEntity(entities.get(0)); + + Assert.notNull(CPEntity); + + try { + + final Batch b = CqlUtils.toDeleteBatchQuery(keyspace.getKeyspace(), tableName, entities, CPEntity, optionsByName); + log.info(b.toString()); + + execute(new SessionCallback() { + + @Override + public Object doInSession(Session s) throws DataAccessException { + + if (deleteAsynchronously) { + s.executeAsync(b); + } else { + s.execute(b); + } + + return null; + + } + }); + + } catch (EntityWriterException e) { + throw exceptionTranslator.translateExceptionIfPossible(new RuntimeException( + "Failed to translate Object to Query", e)); } } + /** + * Insert a row into a Cassandra CQL Table + * + * @param tableName + * @param entity + */ + protected List doBatchInsert(final String tableName, final List entities, + Map optionsByName, final boolean insertAsychronously) { + + Assert.notEmpty(entities); + + CassandraPersistentEntity CPEntity = getEntity(entities.get(0)); + + Assert.notNull(CPEntity); + + try { + + final Batch b = CqlUtils.toInsertBatchQuery(keyspace.getKeyspace(), tableName, entities, CPEntity, optionsByName); + log.info(b.toString()); + + return execute(new SessionCallback>() { + + @Override + public List doInSession(Session s) throws DataAccessException { + + if (insertAsychronously) { + s.executeAsync(b); + } else { + s.execute(b); + } + + return entities; + + } + }); + + } catch (EntityWriterException e) { + throw exceptionTranslator.translateExceptionIfPossible(new RuntimeException( + "Failed to translate Object to Query", e)); + } + } + + /** + * Perform the removal of a Row. + * + * @param tableName + * @param objectToRemove + */ + protected void doDelete(final String tableName, final T objectToRemove, Map optionsByName, + final boolean deleteAsynchronously) { + + CassandraPersistentEntity entity = getEntity(objectToRemove); + + Assert.notNull(entity); + + try { + + final Query q = CqlUtils.toDeleteQuery(keyspace.getKeyspace(), tableName, objectToRemove, entity, optionsByName); + log.info(q.toString()); + + execute(new SessionCallback() { + + @Override + public Object doInSession(Session s) throws DataAccessException { + + if (deleteAsynchronously) { + s.executeAsync(q); + } else { + s.execute(q); + } + + return null; + + } + }); + + } catch (EntityWriterException e) { + throw exceptionTranslator.translateExceptionIfPossible(new RuntimeException( + "Failed to translate Object to Query", e)); + } + } + + /** + * Insert a row into a Cassandra CQL Table + * + * @param tableName + * @param entity + */ + protected T doInsert(final String tableName, final T entity, final Map optionsByName, + final boolean insertAsychronously) { + + CassandraPersistentEntity CPEntity = getEntity(entity); + + Assert.notNull(CPEntity); + + try { + + final Query q = CqlUtils.toInsertQuery(keyspace.getKeyspace(), tableName, entity, CPEntity, optionsByName); + log.info(q.toString()); + + return execute(new SessionCallback() { + + @Override + public T doInSession(Session s) throws DataAccessException { + + if (insertAsychronously) { + s.executeAsync(q); + } else { + s.execute(q); + } + + return entity; + + } + }); + + } catch (EntityWriterException e) { + throw exceptionTranslator.translateExceptionIfPossible(new RuntimeException( + "Failed to translate Object to Query", e)); + } + + } + + /** + * Verify the object is not an iterable type + * + * @param o + */ + protected void ensureNotIterable(Object o) { + if (null != o) { + if (o.getClass().isArray() || ITERABLE_CLASSES.contains(o.getClass().getName())) { + throw new IllegalArgumentException("Cannot use a collection here."); + } + } + } + + /** + * Execute a command at the Session Level + * + * @param callback + * @return + */ + protected T execute(SessionCallback callback) { + + Assert.notNull(callback); + + try { + + return callback.doInSession(session); + + } catch (DataAccessException e) { + throw potentiallyConvertRuntimeException(e); + } + } + + /** + * Determines the PersistentEntityType for a given Object + * + * @param o + * @return + */ + protected CassandraPersistentEntity getEntity(Object o) { + + CassandraPersistentEntity entity = null; + try { + String entityClassName = o.getClass().getName(); + Class entityClass = ClassUtils.forName(entityClassName, beanClassLoader); + entity = mappingContext.getPersistentEntity(entityClass); + } catch (ClassNotFoundException e) { + e.printStackTrace(); + } catch (LinkageError e) { + e.printStackTrace(); + } finally { + } + + return entity; + + } + /** * @param query * @param readRowCallback @@ -305,526 +994,4 @@ public class CassandraTemplate implements CassandraOperations { } } - /** - * @param obj - * @return - */ - private String determineTableName(T obj) { - if (null != obj) { - return determineTableName(obj.getClass()); - } - - return null; - } - - /** - * @param entityClass - * @return - */ - public String determineTableName(Class entityClass) { - - if (entityClass == null) { - throw new InvalidDataAccessApiUsageException( - "No class parameter provided, entity table name can't be determined!"); - } - - CassandraPersistentEntity entity = mappingContext.getPersistentEntity(entityClass); - if (entity == null) { - throw new InvalidDataAccessApiUsageException("No Persitent Entity information found for the class " - + entityClass.getName()); - } - return entity.getTable(); - } - - private RuntimeException potentiallyConvertRuntimeException(RuntimeException ex) { - RuntimeException resolved = this.exceptionTranslator.translateExceptionIfPossible(ex); - return resolved == null ? ex : resolved; - } - - /** - * Insert a row into a Cassandra CQL Table - * - * @param tableName - * @param entity - */ - protected List doBatchInsert(final String tableName, final List entities, final boolean insertAsychronously) { - - Assert.notEmpty(entities); - - CassandraPersistentEntity CPEntity = getEntity(entities.get(0)); - - Assert.notNull(CPEntity); - - try { - - final Batch b = CqlUtils.toInsertBatchQuery(keyspace.getKeyspace(), tableName, entities, CPEntity); - log.info(b.toString()); - - return execute(new SessionCallback>() { - - public List doInSession(Session s) throws DataAccessException { - - if (insertAsychronously) { - s.executeAsync(b); - } else { - s.execute(b); - } - - return entities; - - } - }); - - } catch (EntityWriterException e) { - throw exceptionTranslator.translateExceptionIfPossible(new RuntimeException( - "Failed to translate Object to Query", e)); - } - } - - /** - * Insert a row into a Cassandra CQL Table - * - * @param tableName - * @param entity - */ - protected T doInsert(final String tableName, final T entity, final boolean insertAsychronously) { - - CassandraPersistentEntity CPEntity = getEntity(entity); - - Assert.notNull(CPEntity); - - try { - - final Query q = CqlUtils.toInsertQuery(keyspace.getKeyspace(), tableName, entity, CPEntity); - log.info(q.toString()); - - return execute(new SessionCallback() { - - public T doInSession(Session s) throws DataAccessException { - - if (insertAsychronously) { - s.executeAsync(q); - } else { - s.execute(q); - } - - return entity; - - } - }); - - } catch (EntityWriterException e) { - throw exceptionTranslator.translateExceptionIfPossible(new RuntimeException( - "Failed to translate Object to Query", e)); - } - - } - - /** - * Verify the object is not an iterable type - * - * @param o - */ - protected void ensureNotIterable(Object o) { - if (null != o) { - if (o.getClass().isArray() || ITERABLE_CLASSES.contains(o.getClass().getName())) { - throw new IllegalArgumentException("Cannot use a collection here."); - } - } - } - - /** - * Perform the removal of a Row. - * - * @param objectToRemove - * @param tableName - */ - protected void doDelete(final Object objectToRemove, final String tableName, final boolean deleteAsynchronously) { - - CassandraPersistentEntity entity = getEntity(objectToRemove); - - Assert.notNull(entity); - - try { - - final Query q = CqlUtils.toDeleteQuery(keyspace.getKeyspace(), tableName, objectToRemove, entity); - log.info(q.toString()); - - execute(new SessionCallback() { - - public Object doInSession(Session s) throws DataAccessException { - - if (deleteAsynchronously) { - s.executeAsync(q); - } else { - s.execute(q); - } - - return null; - - } - }); - - } catch (EntityWriterException e) { - throw exceptionTranslator.translateExceptionIfPossible(new RuntimeException( - "Failed to translate Object to Query", e)); - } - } - - /** - * Perform the deletion on a list of objects - * - * @param objectToRemove - * @param tableName - */ - protected void doBatchDelete(final String tableName, final List entities, final boolean deleteAsynchronously) { - - Assert.notEmpty(entities); - - CassandraPersistentEntity CPEntity = getEntity(entities.get(0)); - - Assert.notNull(CPEntity); - - try { - - final Batch b = CqlUtils.toDeleteBatchQuery(keyspace.getKeyspace(), tableName, entities, CPEntity); - log.info(b.toString()); - - execute(new SessionCallback() { - - public Object doInSession(Session s) throws DataAccessException { - - if (deleteAsynchronously) { - s.executeAsync(b); - } else { - s.execute(b); - } - - return null; - - } - }); - - } catch (EntityWriterException e) { - throw exceptionTranslator.translateExceptionIfPossible(new RuntimeException( - "Failed to translate Object to Query", e)); - } - } - - /** - * Execute a command at the Session Level - * - * @param callback - * @return - */ - protected T execute(SessionCallback callback) { - - Assert.notNull(callback); - - try { - - return callback.doInSession(session); - - } catch (DataAccessException e) { - throw potentiallyConvertRuntimeException(e); - } - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#insert(java.lang.Object) - */ - @Override - public T insert(T entity) { - ensureNotIterable(entity); - - String tableName = determineTableName(entity); - - Assert.notNull(tableName); - - return insert(entity, tableName); - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#insert(java.util.List) - */ - @Override - public List insert(List entities) { - - Assert.notNull(entities); - Assert.notEmpty(entities); - - String tableName = getTableName(entities.get(0).getClass()); - - Assert.notNull(tableName); - - return insert(entities, tableName); - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#insert(java.lang.Object, java.lang.String) - */ - @Override - public T insert(T entity, String tableName) { - ensureNotIterable(entity); - return doInsert(tableName, entity, false); - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#insert(java.util.List, java.lang.String) - */ - @Override - public List insert(List entities, String tableName) { - - Assert.notNull(entities); - Assert.notEmpty(entities); - Assert.notNull(tableName); - - return doBatchInsert(tableName, entities, false); - - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#insertAsynchronously(java.lang.Object) - */ - @Override - public T insertAsynchronously(T entity) { - - ensureNotIterable(entity); - - String tableName = determineTableName(entity); - - Assert.notNull(tableName); - - return insertAsynchronously(entity, tableName); - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#insertAsynchronously(java.util.List) - */ - @Override - public List insertAsynchronously(List entities) { - - Assert.notNull(entities); - Assert.notEmpty(entities); - - String tableName = getTableName(entities.get(0).getClass()); - - Assert.notNull(tableName); - - return insertAsynchronously(entities, tableName); - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#insertAsynchronously(java.lang.Object, java.lang.String) - */ - @Override - public T insertAsynchronously(T entity, String tableName) { - - ensureNotIterable(entity); - - return doInsert(tableName, entity, true); - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#insertAsynchronously(java.util.List, java.lang.String) - */ - @Override - public List insertAsynchronously(List entities, String tableName) { - - Assert.notNull(entities); - Assert.notEmpty(entities); - Assert.notNull(tableName); - - return doBatchInsert(tableName, entities, true); - - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#delete(java.lang.Object) - */ - @Override - public void delete(T entity) { - - Assert.notNull(entity); - - String tableName = getTableName(entity.getClass()); - - Assert.notNull(tableName); - - delete(entity, tableName); - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#delete(java.util.List) - */ - @Override - public void delete(List entities) { - - Assert.notNull(entities); - Assert.notEmpty(entities); - - String tableName = getTableName(entities.get(0).getClass()); - - Assert.notNull(tableName); - - delete(entities, tableName); - - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#delete(java.lang.Object, java.lang.String) - */ - @Override - public void delete(T entity, String tableName) { - - CassandraPersistentEntity entityClass = getEntity(entity); - - Assert.notNull(entityClass); - - doDelete(entity, tableName, false); - - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#delete(java.util.List, java.lang.String) - */ - @Override - public void delete(List entities, String tableName) { - - Assert.notNull(entities); - Assert.notEmpty(entities); - Assert.notNull(tableName); - - doBatchDelete(tableName, entities, false); - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#deleteAsychronously(java.lang.Object) - */ - @Override - public void deleteAsychronously(T entity) { - - Assert.notNull(entity); - - String tableName = getTableName(entity.getClass()); - - Assert.notNull(tableName); - - deleteAsychronously(entity, tableName); - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#deleteAsychronously(java.util.List) - */ - @Override - public void deleteAsychronously(List entities) { - - Assert.notNull(entities); - Assert.notEmpty(entities); - - String tableName = getTableName(entities.get(0).getClass()); - - Assert.notNull(tableName); - - deleteAsychronously(entities, tableName); - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#deleteAsychronously(java.lang.Object, java.lang.String) - */ - @Override - public void deleteAsychronously(T entity, String tableName) { - - CassandraPersistentEntity entityClass = getEntity(entity); - - Assert.notNull(entityClass); - - doDelete(entity, tableName, true); - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#deleteAsychronously(java.util.List, java.lang.String) - */ - @Override - public void deleteAsychronously(List entities, String tableName) { - - Assert.notNull(entities); - Assert.notEmpty(entities); - Assert.notNull(tableName); - - doBatchDelete(tableName, entities, true); - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#update(java.lang.Object) - */ - @Override - public T update(T entity) { - // TODO Auto-generated method stub - return null; - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#update(java.util.List) - */ - @Override - public List update(List entities) { - // TODO Auto-generated method stub - return null; - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#update(java.lang.Object, java.lang.String) - */ - @Override - public T update(T entity, String tableName) { - // TODO Auto-generated method stub - return null; - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#update(java.util.List, java.lang.String) - */ - @Override - public List update(List entities, String tableName) { - // TODO Auto-generated method stub - return null; - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#updateAsynchronously(java.lang.Object) - */ - @Override - public T updateAsynchronously(T entity) { - // TODO Auto-generated method stub - return null; - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#updateAsynchronously(java.util.List) - */ - @Override - public List updateAsynchronously(List entities) { - // TODO Auto-generated method stub - return null; - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#updateAsynchronously(java.lang.Object, java.lang.String) - */ - @Override - public T updateAsynchronously(T entity, String tableName) { - // TODO Auto-generated method stub - return null; - } - - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#updateAsynchronously(java.util.List, java.lang.String) - */ - @Override - public List updateAsynchronously(List entities, String tableName) { - // TODO Auto-generated method stub - return null; - } - } diff --git a/src/main/java/org/springframework/data/cassandra/core/ConsistencyLevel.java b/src/main/java/org/springframework/data/cassandra/core/ConsistencyLevel.java new file mode 100644 index 000000000..e8b1247f2 --- /dev/null +++ b/src/main/java/org/springframework/data/cassandra/core/ConsistencyLevel.java @@ -0,0 +1,28 @@ +/* + * Copyright 2011-2013 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.data.cassandra.core; + +/** + * Generic Consistency Levels associated with Cassandra. + * + * @author David Webb + * + */ +public enum ConsistencyLevel { + + ANY, ONE, TWO, THREE, QUOROM, LOCAL_QUOROM, EACH_QUOROM, ALL + +} diff --git a/src/main/java/org/springframework/data/cassandra/core/ConsistencyLevelResolver.java b/src/main/java/org/springframework/data/cassandra/core/ConsistencyLevelResolver.java new file mode 100644 index 000000000..cb02b869c --- /dev/null +++ b/src/main/java/org/springframework/data/cassandra/core/ConsistencyLevelResolver.java @@ -0,0 +1,78 @@ +/* + * Copyright 2011-2013 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.data.cassandra.core; + +/** + * Determine driver consistency level based on ConsistencyLevel + * + * @author David Webb + * + */ +public final class ConsistencyLevelResolver { + + /** + * No instances allowed + */ + private ConsistencyLevelResolver() { + } + + /** + * Decode the generic spring data cassandra enum to the type required by the DataStax Driver. + * + * @param level + * @return The DataStax Driver Consistency Level. + */ + public static com.datastax.driver.core.ConsistencyLevel resolve(ConsistencyLevel level) { + + com.datastax.driver.core.ConsistencyLevel resolvedLevel = com.datastax.driver.core.ConsistencyLevel.ONE; + + /* + * Determine the driver level based on our enum + */ + switch (level) { + case ONE: + resolvedLevel = com.datastax.driver.core.ConsistencyLevel.ONE; + break; + case ALL: + resolvedLevel = com.datastax.driver.core.ConsistencyLevel.ALL; + break; + case ANY: + resolvedLevel = com.datastax.driver.core.ConsistencyLevel.ANY; + break; + case EACH_QUOROM: + resolvedLevel = com.datastax.driver.core.ConsistencyLevel.EACH_QUORUM; + break; + case LOCAL_QUOROM: + resolvedLevel = com.datastax.driver.core.ConsistencyLevel.LOCAL_QUORUM; + break; + case QUOROM: + resolvedLevel = com.datastax.driver.core.ConsistencyLevel.QUORUM; + break; + case THREE: + resolvedLevel = com.datastax.driver.core.ConsistencyLevel.THREE; + break; + case TWO: + resolvedLevel = com.datastax.driver.core.ConsistencyLevel.TWO; + break; + default: + break; + } + + return resolvedLevel; + + } + +} diff --git a/src/main/java/org/springframework/data/cassandra/core/QueryOptions.java b/src/main/java/org/springframework/data/cassandra/core/QueryOptions.java new file mode 100644 index 000000000..86ea87a6b --- /dev/null +++ b/src/main/java/org/springframework/data/cassandra/core/QueryOptions.java @@ -0,0 +1,106 @@ +/* + * Copyright 2011-2013 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.data.cassandra.core; + +import java.util.HashMap; +import java.util.Map; + +/** + * Contains Query Options for Cassnadra queries. This controls the Consistency Tuning and Retry Policy for a Query. + * + * @author David Webb + * + */ +public class QueryOptions { + + private ConsistencyLevel consistencyLevel; + private RetryPolicy retryPolicy; + private Integer ttl; + + /** + * Create a Map of all these options. + */ + public Map toMap() { + + Map m = new HashMap(); + + if (getConsistencyLevel() != null) { + m.put(QueryOptionMapKeys.CONSISTENCY_LEVEL, getConsistencyLevel()); + } + if (getRetryPolicy() != null) { + m.put(QueryOptionMapKeys.RETRY_POLICY, getRetryPolicy()); + } + if (getTtl() != null) { + m.put(QueryOptionMapKeys.TTL, getTtl()); + } + + return m; + } + + /** + * @return Returns the consistencyLevel. + */ + public ConsistencyLevel getConsistencyLevel() { + return consistencyLevel; + } + + /** + * @param consistencyLevel The consistencyLevel to set. + */ + public void setConsistencyLevel(ConsistencyLevel consistencyLevel) { + this.consistencyLevel = consistencyLevel; + } + + /** + * @return Returns the retryPolicy. + */ + public RetryPolicy getRetryPolicy() { + return retryPolicy; + } + + /** + * @param retryPolicy The retryPolicy to set. + */ + public void setRetryPolicy(RetryPolicy retryPolicy) { + this.retryPolicy = retryPolicy; + } + + /** + * @return Returns the ttl. + */ + public Integer getTtl() { + return ttl; + } + + /** + * @param ttl The ttl to set. + */ + public void setTtl(Integer ttl) { + this.ttl = ttl; + } + + /** + * Constants for looking up Map Elements by Key + * + * @author David Webb + * + */ + public static interface QueryOptionMapKeys { + public final String CONSISTENCY_LEVEL = "ConsistencyLevel"; + public final String RETRY_POLICY = "RetryPolicy"; + public final String TTL = "TTL"; + } +} diff --git a/src/main/java/org/springframework/data/cassandra/core/RetryPolicy.java b/src/main/java/org/springframework/data/cassandra/core/RetryPolicy.java new file mode 100644 index 000000000..be617f2e5 --- /dev/null +++ b/src/main/java/org/springframework/data/cassandra/core/RetryPolicy.java @@ -0,0 +1,28 @@ +/* + * Copyright 2011-2013 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.data.cassandra.core; + +/** + * Retry Policies associated with Cassandra. + * + * @author David Webb + * + */ +public enum RetryPolicy { + + DEFAULT, DOWNGRADING_CONSISTENCY, FALLTHROUGH, LOGGING + +} diff --git a/src/main/java/org/springframework/data/cassandra/core/RetryPolicyResolver.java b/src/main/java/org/springframework/data/cassandra/core/RetryPolicyResolver.java new file mode 100644 index 000000000..fbff98cb0 --- /dev/null +++ b/src/main/java/org/springframework/data/cassandra/core/RetryPolicyResolver.java @@ -0,0 +1,67 @@ +/* + * Copyright 2011-2013 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.springframework.data.cassandra.core; + +import com.datastax.driver.core.policies.DefaultRetryPolicy; +import com.datastax.driver.core.policies.DowngradingConsistencyRetryPolicy; +import com.datastax.driver.core.policies.FallthroughRetryPolicy; + +/** + * Determine driver query retry policy + * + * @author David Webb + * + */ +public final class RetryPolicyResolver { + + /** + * No instances allowed + */ + private RetryPolicyResolver() { + } + + /** + * Decode the generic spring data cassandra enum to the type required by the DataStax Driver. + * + * @param level + * @return The DataStax Driver Consistency Level. + */ + public static com.datastax.driver.core.policies.RetryPolicy resolve(RetryPolicy policy) { + + com.datastax.driver.core.policies.RetryPolicy resolvedPolicy = DefaultRetryPolicy.INSTANCE; + + /* + * Determine the driver level based on our enum + */ + switch (policy) { + case DEFAULT: + resolvedPolicy = DefaultRetryPolicy.INSTANCE; + break; + case DOWNGRADING_CONSISTENCY: + resolvedPolicy = DowngradingConsistencyRetryPolicy.INSTANCE; + break; + case FALLTHROUGH: + resolvedPolicy = FallthroughRetryPolicy.INSTANCE; + break; + default: + resolvedPolicy = DefaultRetryPolicy.INSTANCE; + break; + } + + return resolvedPolicy; + + } +} diff --git a/src/main/java/org/springframework/data/cassandra/util/CqlUtils.java b/src/main/java/org/springframework/data/cassandra/util/CqlUtils.java index 5023aafcc..65fc10a1f 100644 --- a/src/main/java/org/springframework/data/cassandra/util/CqlUtils.java +++ b/src/main/java/org/springframework/data/cassandra/util/CqlUtils.java @@ -3,10 +3,16 @@ package org.springframework.data.cassandra.util; import java.lang.reflect.InvocationTargetException; import java.util.ArrayList; import java.util.List; +import java.util.Map; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.dao.InvalidDataAccessApiUsageException; +import org.springframework.data.cassandra.core.ConsistencyLevel; +import org.springframework.data.cassandra.core.ConsistencyLevelResolver; +import org.springframework.data.cassandra.core.QueryOptions; +import org.springframework.data.cassandra.core.RetryPolicy; +import org.springframework.data.cassandra.core.RetryPolicyResolver; import org.springframework.data.cassandra.exception.EntityWriterException; import org.springframework.data.cassandra.mapping.CassandraPersistentEntity; import org.springframework.data.cassandra.mapping.CassandraPersistentProperty; @@ -188,8 +194,6 @@ public abstract class CqlUtils { } }); - // System.out.println("CQL=" + table.asCQLQuery()); - return result; } @@ -200,6 +204,7 @@ public abstract class CqlUtils { * @param tableName * @param entity * @param objectToSave + * @param optionsByName * @param mappingContext * @param beanClassLoader * @@ -207,7 +212,7 @@ public abstract class CqlUtils { * @throws EntityWriterException */ public static Query toInsertQuery(String keyspaceName, String tableName, final Object objectToSave, - CassandraPersistentEntity entity) throws EntityWriterException { + CassandraPersistentEntity entity, Map optionsByName) throws EntityWriterException { final Insert q = QueryBuilder.insertInto(keyspaceName, tableName); final Exception innerException = new Exception(); @@ -242,6 +247,18 @@ public abstract class CqlUtils { throw new EntityWriterException("Failed to convert Persistent Entity to CQL/Query", innerException.getCause()); } + /* + * Add Query Options + */ + addQueryOptions(q, optionsByName); + + /* + * Add TTL to Insert object + */ + if (optionsByName.get(QueryOptions.QueryOptionMapKeys.TTL) != null) { + q.using(QueryBuilder.ttl((Integer) optionsByName.get(QueryOptions.QueryOptionMapKeys.TTL))); + } + return q; } @@ -260,7 +277,8 @@ public abstract class CqlUtils { * @throws EntityWriterException */ public static Batch toInsertBatchQuery(final String keyspaceName, final String tableName, - final List objectsToSave, CassandraPersistentEntity entity) throws EntityWriterException { + final List objectsToSave, CassandraPersistentEntity entity, Map optionsByName) + throws EntityWriterException { /* * Return variable is a Batch statement @@ -271,10 +289,12 @@ public abstract class CqlUtils { for (final T objectToSave : objectsToSave) { - queries.add(toInsertQuery(keyspaceName, tableName, objectToSave, entity)); + queries.add(toInsertQuery(keyspaceName, tableName, objectToSave, entity, optionsByName)); } + addQueryOptions(b, optionsByName); + return b; } @@ -288,7 +308,7 @@ public abstract class CqlUtils { * @throws EntityWriterException */ public static Query toDeleteQuery(String keyspace, String tableName, final Object objectToRemove, - CassandraPersistentEntity entity) throws EntityWriterException { + CassandraPersistentEntity entity, Map optionsByName) throws EntityWriterException { final Delete.Selection ds = QueryBuilder.delete(); final Delete q = ds.from(keyspace, tableName); @@ -328,55 +348,12 @@ public abstract class CqlUtils { throw new EntityWriterException("Failed to convert Persistent Entity to CQL/Query", innerException.getCause()); } + addQueryOptions(q, optionsByName); + return q; } - /** - * Generate the CQL for insert - * - * @param tableName - * @param entity - * @return - */ - public static String toInsertCQL(String tableName, final CassandraPersistentEntity entity) { - - final StringBuilder str = new StringBuilder(); - str.append("INSERT INTO "); - str.append(tableName); - str.append(" ("); - - final List cols = new ArrayList(); - - entity.doWithProperties(new PropertyHandler() { - public void doWithPersistentProperty(CassandraPersistentProperty prop) { - - if (str.charAt(str.length() - 1) != '(') { - str.append(", "); - } - - String columnName = prop.getColumnName(); - cols.add(columnName); - - str.append(columnName); - - } - }); - - str.append(") VALUES ("); - - for (int i = 0; i < cols.size(); i++) { - if (i > 0) { - str.append(", "); - } - str.append("?"); - } - - str.append(")"); - - return str.toString(); - } - /** * @param dataType * @return @@ -423,7 +400,7 @@ public abstract class CqlUtils { * @throws EntityWriterException */ public static Batch toDeleteBatchQuery(String keyspaceName, String tableName, List entities, - CassandraPersistentEntity entity) throws EntityWriterException { + CassandraPersistentEntity entity, Map optionsByName) throws EntityWriterException { /* * Return variable is a Batch statement @@ -434,12 +411,40 @@ public abstract class CqlUtils { for (final T objectToSave : entities) { - queries.add(toDeleteQuery(keyspaceName, tableName, objectToSave, entity)); + queries.add(toDeleteQuery(keyspaceName, tableName, objectToSave, entity, optionsByName)); } + addQueryOptions(b, optionsByName); + return b; } + /** + * Add common Query options for all types of queries. + * + * @param q + * @param optionsByName + */ + private static void addQueryOptions(Query q, Map optionsByName) { + + if (optionsByName == null) { + return; + } + + /* + * Add Query Options + */ + if (optionsByName.get(QueryOptions.QueryOptionMapKeys.CONSISTENCY_LEVEL) != null) { + q.setConsistencyLevel(ConsistencyLevelResolver.resolve((ConsistencyLevel) optionsByName + .get(QueryOptions.QueryOptionMapKeys.CONSISTENCY_LEVEL))); + } + if (optionsByName.get(QueryOptions.QueryOptionMapKeys.RETRY_POLICY) != null) { + q.setRetryPolicy(RetryPolicyResolver.resolve((RetryPolicy) optionsByName + .get(QueryOptions.QueryOptionMapKeys.RETRY_POLICY))); + } + + } + }