From 471fae991160c4776064869edc5f1cbebc2e2065 Mon Sep 17 00:00:00 2001 From: Matthew Adams Date: Thu, 6 Feb 2014 12:04:35 -0600 Subject: [PATCH] DATACASS-33 - simplifying CassandraOperations --- .../core/AsynchronousQueryListener.java | 7 +- .../cassandra/core/CqlOperations.java | 33 +- .../cassandra/core/CqlTemplate.java | 42 +- .../cassandra/core/util/CollectionUtils.java | 5 + .../convert/MappingCassandraConverter.java | 5 + .../core/CassandraAdminTemplate.java | 20 - ...ava => CassandraConverterRowCallback.java} | 18 +- .../cassandra/core/CassandraOperations.java | 231 +----- .../cassandra/core/CassandraTemplate.java | 716 ++++-------------- .../mapping/CassandraMappingContext.java | 21 + .../DefaultCassandraMappingContext.java | 25 + .../support/SimpleCassandraRepository.java | 9 +- .../template/CassandraDataOperationsTest.java | 166 ++-- 13 files changed, 395 insertions(+), 903 deletions(-) rename spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/{ReadRowCallback.java => CassandraConverterRowCallback.java} (68%) diff --git a/spring-cassandra/src/main/java/org/springframework/cassandra/core/AsynchronousQueryListener.java b/spring-cassandra/src/main/java/org/springframework/cassandra/core/AsynchronousQueryListener.java index 3c59206b4..7cf4a587f 100644 --- a/spring-cassandra/src/main/java/org/springframework/cassandra/core/AsynchronousQueryListener.java +++ b/spring-cassandra/src/main/java/org/springframework/cassandra/core/AsynchronousQueryListener.java @@ -18,16 +18,17 @@ package org.springframework.cassandra.core; import com.datastax.driver.core.ResultSetFuture; /** - * @author David Webb + * Interface used to give an implementation access to a {@link ResultSetFuture} after the query has completed. * + * @author David Webb */ public interface AsynchronousQueryListener { /** * Called upon Query Completion. * - * @param rsf The given ResultSetFuture's get methods should return immediately. + * @param rsf The {@link ResultSetFuture}. Since this isn't called until the asynchronous query completes, it can be + * immediately interrogated. */ public void onQueryComplete(ResultSetFuture rsf); - } diff --git a/spring-cassandra/src/main/java/org/springframework/cassandra/core/CqlOperations.java b/spring-cassandra/src/main/java/org/springframework/cassandra/core/CqlOperations.java index fa9b5423b..66842d696 100644 --- a/spring-cassandra/src/main/java/org/springframework/cassandra/core/CqlOperations.java +++ b/spring-cassandra/src/main/java/org/springframework/cassandra/core/CqlOperations.java @@ -56,6 +56,7 @@ public interface CqlOperations { * Executes the supplied CQL Query and returns nothing. * * @param cql + * @see #query(String) */ void execute(String cql) throws DataAccessException; @@ -108,10 +109,10 @@ public interface CqlOperations { ResultSetFuture queryAsynchronously(String cql, QueryOptions options); /** - * Executes the provided CQL Query with the provided Runnable implementations. + * Executes the provided CQL Query with the provided {@link Runnable} implementation. * * @param cql The Query - * @param listener Runnable Listener for handling the query in a separate thread + * @param listener {@link Runnable} listener for handling the query in a separate thread */ void queryAsynchronously(String cql, Runnable listener); @@ -121,7 +122,8 @@ public interface CqlOperations { * query is completed for optimal flexibility. * * @param cql The Query - * @param listener Runnable Listener for handling the query in a separate thread + * @param listener {@link AsynchronousQueryListener} Listener for handling the query's {@link ResultSetFuture} in a + * separate thread */ void queryAsynchronously(String cql, AsynchronousQueryListener listener); @@ -189,6 +191,23 @@ public interface CqlOperations { */ void queryAsynchronously(String cql, AsynchronousQueryListener listener, QueryOptions options, Executor executor); + /** + * Executes the provided CQL query and returns the {@link ResultSet}. + * + * @param cql The query + * @return The {@link ResultSet} + */ + ResultSet query(String cql); + + /** + * Executes the provided CQL query with the given {@link QueryOptions} and returns the {@link ResultSet}. + * + * @param cql The query + * @param options The {@link QueryOptions}; may be null. + * @return The {@link ResultSet} + */ + ResultSet query(String cql, QueryOptions options); + /** * Executes the provided CQL Query, and extracts the results with the ResultSetExtractor. * @@ -762,6 +781,14 @@ public interface CqlOperations { */ void truncate(String tableName); + /** + * Counts all rows for given table + * + * @param tableName + * @return + */ + long count(String tableName); + /** * Convenience method to convert the given specification to CQL and execute it. * diff --git a/spring-cassandra/src/main/java/org/springframework/cassandra/core/CqlTemplate.java b/spring-cassandra/src/main/java/org/springframework/cassandra/core/CqlTemplate.java index 52e8b034e..b7a4363e3 100644 --- a/spring-cassandra/src/main/java/org/springframework/cassandra/core/CqlTemplate.java +++ b/spring-cassandra/src/main/java/org/springframework/cassandra/core/CqlTemplate.java @@ -45,6 +45,7 @@ import org.springframework.cassandra.core.keyspace.DropKeyspaceSpecification; import org.springframework.cassandra.core.keyspace.DropTableSpecification; import org.springframework.cassandra.support.CassandraAccessor; import org.springframework.dao.DataAccessException; +import org.springframework.dao.InvalidDataAccessApiUsageException; import org.springframework.dao.QueryTimeoutException; import org.springframework.util.Assert; @@ -63,6 +64,7 @@ import com.datastax.driver.core.SimpleStatement; import com.datastax.driver.core.Statement; import com.datastax.driver.core.exceptions.DriverException; import com.datastax.driver.core.querybuilder.QueryBuilder; +import com.datastax.driver.core.querybuilder.Select; import com.datastax.driver.core.querybuilder.Truncate; /** @@ -304,6 +306,23 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { return process(doExecute(cql, options), rowMapper); } + @Override + public ResultSet query(String cql) { + return query(cql, (QueryOptions) null); + } + + @Override + public ResultSet query(String cql, QueryOptions options) { + + return query(cql, new ResultSetExtractor() { + + @Override + public ResultSet extractData(ResultSet rs) throws DriverException, DataAccessException { + return rs; + } + }, options); + } + @Override public List query(String cql, RowMapper rowMapper) throws DataAccessException { return query(cql, rowMapper, null); @@ -434,7 +453,7 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { * * @return */ - private Set getHosts() { + protected Set getHosts() { /* * Get the cluster metadata for this session @@ -965,4 +984,25 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations { }); } + @Override + public long count(String tableName) { + return selectCount(QueryBuilder.select().countAll().from(tableName).getQueryString()); + } + + protected long selectCount(String countQuery) { + + return query(countQuery, new ResultSetExtractor() { + + @Override + public Long extractData(ResultSet rs) throws DriverException, DataAccessException { + + Row row = rs.one(); + if (row == null) { + throw new InvalidDataAccessApiUsageException(String.format("count query did not return any results")); + } + + return row.getLong(0); + } + }); + } } \ No newline at end of file diff --git a/spring-cassandra/src/main/java/org/springframework/cassandra/core/util/CollectionUtils.java b/spring-cassandra/src/main/java/org/springframework/cassandra/core/util/CollectionUtils.java index 8004ca8e9..e9ceb2b9a 100644 --- a/spring-cassandra/src/main/java/org/springframework/cassandra/core/util/CollectionUtils.java +++ b/spring-cassandra/src/main/java/org/springframework/cassandra/core/util/CollectionUtils.java @@ -5,6 +5,11 @@ import java.util.List; public class CollectionUtils { + @SuppressWarnings("unchecked") + public static T[] toArray(Iterable i) { + return (T[]) toList(i).toArray(); + } + public static List toList(Iterable i) { List list = null; diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/convert/MappingCassandraConverter.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/convert/MappingCassandraConverter.java index 8315f0b19..ec64f8c90 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/convert/MappingCassandraConverter.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/convert/MappingCassandraConverter.java @@ -34,6 +34,7 @@ import org.springframework.data.mapping.model.MappingException; import org.springframework.data.mapping.model.SpELContext; import org.springframework.data.util.ClassTypeInformation; import org.springframework.data.util.TypeInformation; +import org.springframework.util.Assert; import org.springframework.util.ClassUtils; import com.datastax.driver.core.Row; @@ -67,7 +68,11 @@ public class MappingCassandraConverter extends AbstractCassandraConverter implem * @param mappingContext must not be {@literal null}. */ public MappingCassandraConverter(CassandraMappingContext mappingContext) { + super(new DefaultConversionService()); + + Assert.notNull(mappingContext); + this.mappingContext = mappingContext; this.spELContext = new SpELContext(RowReaderPropertyAccessor.INSTANCE); } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraAdminTemplate.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraAdminTemplate.java index 42ace67b9..c7641f026 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraAdminTemplate.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraAdminTemplate.java @@ -124,24 +124,4 @@ public class CassandraAdminTemplate extends CassandraTemplate implements Cassand } }); } - - /** - * @param entityClass - * @return - */ - @Override - public String determineTableName(Class entityClass) { - - if (entityClass == null) { - throw new InvalidDataAccessApiUsageException( - "No class parameter provided, entity table name can't be determined!"); - } - - CassandraPersistentEntity entity = getCassandraMappingContext().getPersistentEntity(entityClass); - if (entity == null) { - throw new InvalidDataAccessApiUsageException("No Persitent Entity information found for the class " - + entityClass.getName()); - } - return entity.getTableName(); - } } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReadRowCallback.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraConverterRowCallback.java similarity index 68% rename from spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReadRowCallback.java rename to spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraConverterRowCallback.java index a87050ebb..8b02200de 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/ReadRowCallback.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraConverterRowCallback.java @@ -16,32 +16,34 @@ package org.springframework.data.cassandra.core; import org.springframework.cassandra.core.RowCallback; -import org.springframework.data.convert.EntityReader; +import org.springframework.data.cassandra.convert.CassandraConverter; import org.springframework.util.Assert; import com.datastax.driver.core.Row; /** - * Simple {@link RowCallback} that will transform {@link Row} into the given target type using the given - * {@link EntityReader}. + * Simple {@link RowCallback} that will transform a {@link Row} into the given target type using the given + * {@link CassandraConverter}. * * @author Alex Shvid + * @author Matthew T. Adams */ -public class ReadRowCallback implements RowCallback { +public class CassandraConverterRowCallback implements RowCallback { - private final EntityReader reader; + private final CassandraConverter reader; private final Class type; - public ReadRowCallback(EntityReader reader, Class type) { + public CassandraConverterRowCallback(CassandraConverter reader, Class type) { + Assert.notNull(reader); Assert.notNull(type); + this.reader = reader; this.type = type; } @Override public T doWith(Row row) { - T source = reader.read(type, row); - return source; + return reader.read(type, row); } } \ No newline at end of file diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraOperations.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraOperations.java index fcbe05844..e0e5d8586 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraOperations.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraOperations.java @@ -21,8 +21,6 @@ import org.springframework.cassandra.core.CqlOperations; import org.springframework.cassandra.core.QueryOptions; import org.springframework.data.cassandra.convert.CassandraConverter; -import com.datastax.driver.core.querybuilder.Select; - /** * Operations for interacting with Cassandra. These operations are used by the Repository implementation, but can also * be used directly when that is desired by the developer. @@ -45,50 +43,25 @@ public interface CassandraOperations extends CqlOperations { * Execute query and convert ResultSet to the list of entities * * @param query must not be {@literal null}. - * @param selectClass must not be {@literal null}, mapped entity type. + * @param type must not be {@literal null}, mapped entity type. * @return */ - List select(String cql, Class selectClass); + List select(String cql, Class type); - /** - * Execute query and convert ResultSet to the list of entities - * - * @param selectQuery must not be {@literal null}. - * @param selectClass must not be {@literal null}, mapped entity type. - * @return - */ - List select(Select selectQuery, Class selectClass); - - T selectOneById(Class selectClass, Object id); + T selectOneById(Class type, Object id); /** * Execute query and convert ResultSet to the entity * * @param query must not be {@literal null}. - * @param selectClass must not be {@literal null}, mapped entity type. + * @param type must not be {@literal null}, mapped entity type. * @return */ - T selectOne(String cql, Class selectClass); + T selectOne(String cql, Class type); - T selectOne(Select selectQuery, Class selectClass); + boolean exists(Class type, Object id); - Long countById(Class clazz, Object id); - - /** - * Counts rows for given query - * - * @param selectQuery - * @return - */ - Long count(Select selectQuery); - - /** - * Counts all rows for given table - * - * @param tableName - * @return - */ - Long count(String tableName); + long count(Class type); /** * Insert the given object to the table by id. @@ -97,23 +70,6 @@ public interface CassandraOperations extends CqlOperations { */ T insert(T entity); - /** - * Insert the given object to the table by id. - * - * @param entity - * @param tableName - * @return - */ - 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 @@ -130,15 +86,6 @@ public interface CassandraOperations extends CqlOperations { */ List insert(List entities); - /** - * Insert the given list of objects to the table by name. - * - * @param entities - * @param tableName - * @return - */ - List insert(List entities, String tableName); - /** * @param entities * @param tableName @@ -147,14 +94,6 @@ public interface CassandraOperations extends CqlOperations { */ List insert(List entities, QueryOptions options); - /** - * @param entities - * @param tableName - * @param options - * @return - */ - List insert(List entities, String tableName, QueryOptions options); - /** * Insert the given object to the table by id. * @@ -162,13 +101,6 @@ public interface CassandraOperations extends CqlOperations { */ 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 @@ -177,14 +109,6 @@ public interface CassandraOperations extends CqlOperations { */ T insertAsynchronously(T entity, QueryOptions options); - /** - * @param entity - * @param tableName - * @param options - * @return - */ - T insertAsynchronously(T entity, String tableName, QueryOptions options); - /** * Insert the given object to the table by id. * @@ -192,13 +116,6 @@ public interface CassandraOperations extends CqlOperations { */ List insertAsynchronously(List entities); - /** - * Insert the given object to the table by id. - * - * @param object - */ - List insertAsynchronously(List entities, String tableName); - /** * @param entities * @param tableName @@ -207,14 +124,6 @@ public interface CassandraOperations extends CqlOperations { */ List insertAsynchronously(List entities, QueryOptions options); - /** - * @param entities - * @param tableName - * @param options - * @return - */ - List insertAsynchronously(List entities, String tableName, QueryOptions options); - /** * Insert the given object to the table by id. * @@ -222,13 +131,6 @@ public interface CassandraOperations extends CqlOperations { */ 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 @@ -237,14 +139,6 @@ public interface CassandraOperations extends CqlOperations { */ T update(T entity, QueryOptions options); - /** - * @param entity - * @param tableName - * @param options - * @return - */ - T update(T entity, String tableName, QueryOptions options); - /** * Insert the given object to the table by id. * @@ -252,13 +146,6 @@ public interface CassandraOperations extends CqlOperations { */ List update(List entities); - /** - * Insert the given object to the table by id. - * - * @param object - */ - List update(List entities, String tableName); - /** * @param entities * @param tableName @@ -267,14 +154,6 @@ public interface CassandraOperations extends CqlOperations { */ List update(List entities, QueryOptions options); - /** - * @param entities - * @param tableName - * @param options - * @return - */ - List update(List entities, String tableName, QueryOptions options); - /** * Insert the given object to the table by id. * @@ -282,13 +161,6 @@ public interface CassandraOperations extends CqlOperations { */ T updateAsynchronously(T entity); - /** - * Insert the given object to the table by id. - * - * @param object - */ - T updateAsynchronously(T entity, String tableName); - /** * @param entity * @param tableName @@ -297,14 +169,6 @@ public interface CassandraOperations extends CqlOperations { */ T updateAsynchronously(T entity, QueryOptions options); - /** - * @param entity - * @param tableName - * @param options - * @return - */ - T updateAsynchronously(T entity, String tableName, QueryOptions options); - /** * Insert the given object to the table by id. * @@ -312,13 +176,6 @@ public interface CassandraOperations extends CqlOperations { */ List updateAsynchronously(List entities); - /** - * Insert the given object to the table by id. - * - * @param object - */ - List updateAsynchronously(List entities, String tableName); - /** * @param entities * @param tableName @@ -327,14 +184,6 @@ public interface CassandraOperations extends CqlOperations { */ List updateAsynchronously(List entities, QueryOptions options); - /** - * @param entities - * @param tableName - * @param options - * @return - */ - List updateAsynchronously(List entities, String tableName, QueryOptions options); - /** * Remove the given object from the table by id. * @@ -342,14 +191,6 @@ public interface CassandraOperations extends CqlOperations { */ 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 @@ -357,13 +198,6 @@ public interface CassandraOperations extends CqlOperations { */ void delete(T entity, QueryOptions options); - /** - * @param entity - * @param tableName - * @param options - */ - void delete(T entity, String tableName, QueryOptions options); - /** * Remove the given object from the table by id. * @@ -371,14 +205,6 @@ public interface CassandraOperations extends CqlOperations { */ void delete(List entities); - /** - * Removes the given object from the given table. - * - * @param object - * @param table must not be {@literal null} or empty. - */ - void delete(List entities, String tableName); - /** * @param entities * @param tableName @@ -386,13 +212,6 @@ public interface CassandraOperations extends CqlOperations { */ void delete(List entities, QueryOptions options); - /** - * @param entities - * @param tableName - * @param options - */ - void delete(List entities, String tableName, QueryOptions options); - /** * Remove the given object from the table by id. * @@ -407,21 +226,6 @@ public interface CassandraOperations extends CqlOperations { */ void deleteAsynchronously(T entity, QueryOptions options); - /** - * @param entity - * @param tableName - * @param options - */ - void deleteAsynchronously(T entity, String tableName, QueryOptions options); - - /** - * Removes the given object from the given table. - * - * @param object - * @param table must not be {@literal null} or empty. - */ - void deleteAsynchronously(T entity, String tableName); - /** * Remove the given object from the table by id. * @@ -429,14 +233,6 @@ public interface CassandraOperations extends CqlOperations { */ void deleteAsynchronously(List entities); - /** - * Removes the given object from the given table. - * - * @param object - * @param table must not be {@literal null} or empty. - */ - void deleteAsynchronously(List entities, String tableName); - /** * @param entities * @param tableName @@ -444,13 +240,6 @@ public interface CassandraOperations extends CqlOperations { */ void deleteAsynchronously(List entities, QueryOptions options); - /** - * @param entities - * @param tableName - * @param options - */ - void deleteAsynchronously(List entities, String tableName, QueryOptions options); - /** * Returns the underlying {@link CassandraConverter}. * @@ -458,9 +247,9 @@ public interface CassandraOperations extends CqlOperations { */ CassandraConverter getConverter(); - void deleteById(Class clazz, Object id); + void deleteById(Class type, Object id); - List selectByIds(Class clazz, Iterable ids); + List selectByIds(Class type, Iterable ids); - List selectAll(Class clazz); + List selectAll(Class type); } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java index 2bfc6e9bb..6f5ce78c9 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java @@ -16,16 +16,15 @@ package org.springframework.data.cassandra.core; import java.util.ArrayList; -import java.util.Collections; import java.util.Iterator; import java.util.List; import org.springframework.cassandra.core.CqlTemplate; import org.springframework.cassandra.core.QueryOptions; import org.springframework.cassandra.core.SessionCallback; +import org.springframework.cassandra.core.util.CollectionUtils; import org.springframework.dao.DataAccessException; import org.springframework.dao.DuplicateKeyException; -import org.springframework.dao.InvalidDataAccessApiUsageException; import org.springframework.data.cassandra.convert.CassandraConverter; import org.springframework.data.cassandra.mapping.CassandraMappingContext; import org.springframework.data.cassandra.mapping.CassandraPersistentEntity; @@ -34,13 +33,10 @@ import org.springframework.data.convert.EntityWriter; import org.springframework.data.mapping.PropertyHandler; import org.springframework.data.mapping.model.BeanWrapper; import org.springframework.util.Assert; -import org.springframework.util.CollectionUtils; -import com.datastax.driver.core.Query; import com.datastax.driver.core.ResultSet; import com.datastax.driver.core.Row; import com.datastax.driver.core.Session; -import com.datastax.driver.core.Statement; import com.datastax.driver.core.querybuilder.Batch; import com.datastax.driver.core.querybuilder.Clause; import com.datastax.driver.core.querybuilder.Delete; @@ -62,13 +58,9 @@ import com.datastax.driver.core.querybuilder.Update; */ public class CassandraTemplate extends CqlTemplate implements CassandraOperations { - /* - * Required elements for successful Template Operations. These can be set with the Constructor, or wired in - * later. - */ - private CassandraConverter cassandraConverter; - private CassandraMappingContext mappingContext; - private boolean useFieldAccessOnly = false; + protected CassandraConverter cassandraConverter; + protected CassandraMappingContext mappingContext; + protected boolean useFieldAccessOnly = false; /** * Default Constructor for wiring in the required components later @@ -90,9 +82,9 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation public void setConverter(CassandraConverter cassandraConverter) { Assert.notNull(cassandraConverter); + this.cassandraConverter = cassandraConverter; mappingContext = cassandraConverter.getCassandraMappingContext(); - Assert.notNull(mappingContext); } @Override @@ -104,6 +96,19 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation return mappingContext; } + public boolean getUseFieldAccessOnly() { + return useFieldAccessOnly; + } + + /** + * Whether only fields should be used when accessing a persistent entity's data. + * + * @param useFieldAccessOnly + */ + public void setUseFieldAccessOnly(boolean useFieldAccessOnly) { + this.useFieldAccessOnly = useFieldAccessOnly; + } + @Override public void afterPropertiesSet() { super.afterPropertiesSet(); @@ -113,71 +118,41 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation } @Override - public Long countById(Class clazz, Object id) { + public boolean exists(Class type, Object id) { - Assert.notNull(clazz); + Assert.notNull(type); Assert.notNull(id); - CassandraPersistentEntity entity = mappingContext.getPersistentEntity(clazz); - if (entity == null) { - throw new IllegalArgumentException(String.format("unknown persistent class [%s]", clazz.getName())); - } + CassandraPersistentEntity entity = mappingContext.getRequiredPersistentEntity(type); Select select = QueryBuilder.select().countAll().from(entity.getTableName()); appendIdCriteria(select.where(), entity, id); - return count(select); + return count(select.getQueryString()) != 0; } @Override - public Long count(Select selectQuery) { - return selectCount(selectQuery); - } - - @Override - public Long count(String tableName) { - Select select = QueryBuilder.select().countAll().from(tableName); - return selectCount(select); + public long count(Class type) { + return count(getTableName(type)); } @Override public void delete(List entities) { - - Assert.notEmpty(entities); - - String tableName = getTableName(entities.get(0).getClass()); - Assert.notNull(tableName); - - delete(entities, tableName); + delete(entities, null); } @Override public void delete(List entities, QueryOptions options) { - String tableName = getTableName(entities.get(0).getClass()); - Assert.notNull(tableName); - delete(entities, tableName, options); + batchDelete(entities, options, false); } @Override - public void delete(List entities, String tableName) { - delete(entities, tableName, null); - } + public void deleteById(Class type, Object id) { - @Override - public void delete(List entities, String tableName, QueryOptions options) { - Assert.notNull(entities); - Assert.notEmpty(entities); - Assert.notNull(tableName); - batchDelete(tableName, entities, options, false); - } - - @Override - public void deleteById(Class clazz, Object id) { - - Assert.notNull(clazz); + Assert.notNull(type); Assert.notNull(id); - CassandraPersistentEntity entity = mappingContext.getPersistentEntity(clazz); + CassandraPersistentEntity entity = mappingContext.getRequiredPersistentEntity(type); Delete delete = QueryBuilder.delete().all().from(entity.getTableName()); appendIdCriteria(delete.where(), entity, id); @@ -187,295 +162,129 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation @Override public void delete(T entity) { - String tableName = getTableName(entity.getClass()); - Assert.notNull(tableName); - delete(entity, tableName); + delete(entity, null); } @Override public void delete(T entity, QueryOptions options) { - String tableName = getTableName(entity.getClass()); - Assert.notNull(tableName); - delete(entity, tableName, options); - } - - @Override - public void delete(T entity, String tableName) { - delete(entity, tableName, null); - } - - @Override - public void delete(T entity, String tableName, QueryOptions options) { - Assert.notNull(entity); - Assert.notNull(tableName); - delete(tableName, entity, options, false); + delete(entity, options, false); } @Override public void deleteAsynchronously(List entities) { - String tableName = getTableName(entities.get(0).getClass()); - Assert.notNull(tableName); - deleteAsynchronously(entities, tableName); + deleteAsynchronously(entities, null); } @Override public void deleteAsynchronously(List entities, QueryOptions options) { - String tableName = getTableName(entities.get(0).getClass()); - Assert.notNull(tableName); - deleteAsynchronously(entities, tableName, options); - } - - @Override - public void deleteAsynchronously(List entities, String tableName) { - deleteAsynchronously(entities, tableName, null); - } - - @Override - public void deleteAsynchronously(List entities, String tableName, QueryOptions options) { - Assert.notNull(entities); - Assert.notEmpty(entities); - Assert.notNull(tableName); - batchDelete(tableName, entities, options, true); + batchDelete(entities, options, true); } @Override public void deleteAsynchronously(T entity) { - String tableName = getTableName(entity.getClass()); - Assert.notNull(tableName); - deleteAsynchronously(entity, tableName); + deleteAsynchronously(entity, null); } @Override public void deleteAsynchronously(T entity, QueryOptions options) { - String tableName = getTableName(entity.getClass()); - Assert.notNull(tableName); - deleteAsynchronously(entity, tableName, options); + delete(entity, options, true); } @Override - public void deleteAsynchronously(T entity, String tableName) { - deleteAsynchronously(entity, tableName, null); - } - - @Override - public void deleteAsynchronously(T entity, String tableName, QueryOptions options) { - Assert.notNull(entity); - Assert.notNull(tableName); - delete(tableName, entity, options, true); - } - - /** - * @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.getTableName(); - } - - @Override - public String getTableName(Class entityClass) { - return determineTableName(entityClass); + public String getTableName(Class type) { + return mappingContext.getRequiredPersistentEntity(type).getTableName(); } @Override public List insert(List entities) { - String tableName = getTableName(entities.get(0).getClass()); - Assert.notNull(tableName); - return insert(entities, tableName); + return insert(entities, null); } @Override public List insert(List entities, QueryOptions options) { - String tableName = getTableName(entities.get(0).getClass()); - Assert.notNull(tableName); - return insert(entities, tableName, options); - } - - @Override - public List insert(List entities, String tableName) { - return insert(entities, tableName, null); - } - - @Override - public List insert(List entities, String tableName, QueryOptions options) { - Assert.notNull(entities); - Assert.notEmpty(entities); - Assert.notNull(tableName); - return batchInsert(tableName, entities, options, false); + return batchInsert(entities, options, false); } @Override public T insert(T entity) { - String tableName = determineTableName(entity); - Assert.notNull(tableName); - return insert(entity, tableName); + return insert(entity, null); } @Override public T insert(T entity, QueryOptions options) { - String tableName = determineTableName(entity); - Assert.notNull(tableName); - return insert(entity, tableName, options); - } - - @Override - public T insert(T entity, String tableName) { - return insert(entity, tableName, null); - } - - @Override - public T insert(T entity, String tableName, QueryOptions options) { - Assert.notNull(entity); - Assert.notNull(tableName); - ensureNotIterable(entity); - return insert(tableName, entity, options, false); + return insert(entity, options, false); } @Override public List insertAsynchronously(List entities) { - String tableName = getTableName(entities.get(0).getClass()); - Assert.notNull(tableName); - return insertAsynchronously(entities, tableName); + return insertAsynchronously(entities, null); } @Override public List insertAsynchronously(List entities, QueryOptions options) { - String tableName = getTableName(entities.get(0).getClass()); - Assert.notNull(tableName); - return insertAsynchronously(entities, tableName, options); - } - - @Override - public List insertAsynchronously(List entities, String tableName) { - return insertAsynchronously(entities, tableName, null); - } - - @Override - public List insertAsynchronously(List entities, String tableName, QueryOptions options) { - Assert.notNull(entities); - Assert.notEmpty(entities); - Assert.notNull(tableName); - return batchInsert(tableName, entities, options, true); + return batchInsert(entities, options, true); } @Override public T insertAsynchronously(T entity) { - String tableName = determineTableName(entity); - Assert.notNull(tableName); - return insertAsynchronously(entity, tableName); + return insertAsynchronously(entity, null); } @Override public T insertAsynchronously(T entity, QueryOptions options) { - String tableName = determineTableName(entity); - Assert.notNull(tableName); - return insertAsynchronously(entity, tableName, options); + return insert(entity, options, true); } @Override - public T insertAsynchronously(T entity, String tableName) { - return insertAsynchronously(entity, tableName, null); + public List selectAll(Class type) { + return select(QueryBuilder.select().all().from(getTableName(type)).getQueryString(), type); } @Override - public T insertAsynchronously(T entity, String tableName, QueryOptions options) { - Assert.notNull(entity); - Assert.notNull(tableName); - - ensureNotIterable(entity); - - return insert(tableName, entity, options, true); - } - - @Override - public List selectAll(Class selectClass) { - - Assert.notNull(selectClass); - - CassandraPersistentEntity entity = mappingContext.getPersistentEntity(selectClass); - if (entity == null) { - throw new IllegalArgumentException(String.format("unknown persistent class [%s]", selectClass.getName())); - } - return select(QueryBuilder.select().all().from(entity.getTableName()), selectClass); - } - - @Override - public List select(Select cql, Class selectClass) { - - Assert.notNull(cql); - - return select(cql.getQueryString(), selectClass); - } - - @Override - public List select(String cql, Class selectClass) { + public List select(String cql, Class type) { Assert.hasText(cql); - Assert.notNull(selectClass); + Assert.notNull(type); - return select(cql, new ReadRowCallback(cassandraConverter, selectClass)); + return select(cql, new CassandraConverterRowCallback(cassandraConverter, type)); } - @SuppressWarnings({ "rawtypes", "unchecked" }) @Override - public List selectByIds(Class clazz, Iterable ids) { + public List selectByIds(Class type, Iterable ids) { + + CassandraPersistentEntity entity = mappingContext.getRequiredPersistentEntity(type); - CassandraPersistentEntity entity = mappingContext.getPersistentEntity(clazz); - if (entity == null) { - throw new IllegalArgumentException(String.format("unknown persistent entity class [%s]", clazz.getName())); - } if (entity.getIdProperty().isCompositePrimaryKey()) { throw new IllegalArgumentException(String.format( - "entity class [%s] uses a composite primary key class [%s] which this method can't support", clazz.getName(), + "entity class [%s] uses a composite primary key class [%s] which this method can't support", type.getName(), entity.getIdProperty().getCompositePrimaryKeyEntity().getType().getName())); } - List idList = null; - if (ids instanceof List) { - idList = (List) ids; - } else { - idList = new ArrayList(); - for (Object id : ids) { - idList.add(id); - } - } - Select select = QueryBuilder.select().all().from(entity.getTableName()); - select.where(QueryBuilder.in(entity.getIdProperty().getColumnName(), idList.toArray())); + select.where(QueryBuilder.in(entity.getIdProperty().getColumnName(), CollectionUtils.toArray(ids))); - return select(select, clazz); + return select(select.getQueryString(), type); } @Override - public T selectOneById(Class selectClass, Object id) { + public T selectOneById(Class type, Object id) { - Assert.notNull(selectClass); + Assert.notNull(type); Assert.notNull(id); - CassandraPersistentEntity entityClass = mappingContext.getPersistentEntity(selectClass); - if (entityClass == null) { - throw new IllegalArgumentException(String.format("unknown entity class [%s]", selectClass.getName())); + CassandraPersistentEntity entity = mappingContext.getPersistentEntity(type); + if (entity == null) { + throw new IllegalArgumentException(String.format("unknown entity class [%s]", type.getName())); } - Select select = QueryBuilder.select().all().from(entityClass.getTableName()); - appendIdCriteria(select.where(), entityClass, id); + Select select = QueryBuilder.select().all().from(entity.getTableName()); + appendIdCriteria(select.where(), entity, id); - return selectOne(select, selectClass); + return selectOne(select.getQueryString(), type); } protected interface ClauseCallback { - void onClause(Clause clause); + void doWithClause(Clause clause); } protected void appendIdCriteria(final ClauseCallback clauseCallback, CassandraPersistentEntity entity, Object id) { @@ -494,7 +303,7 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation @Override public void doWithPersistentProperty(CassandraPersistentProperty p) { - clauseCallback.onClause(QueryBuilder.eq(p.getColumnName(), + clauseCallback.doWithClause(QueryBuilder.eq(p.getColumnName(), idWrapper.getProperty(p, p.getActualType(), useFieldAccessOnly))); } }); @@ -502,7 +311,7 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation return; } - clauseCallback.onClause(QueryBuilder.eq(idProperty.getColumnName(), id)); + clauseCallback.doWithClause(QueryBuilder.eq(idProperty.getColumnName(), id)); } protected void appendIdCriteria(final com.datastax.driver.core.querybuilder.Select.Where where, @@ -511,7 +320,7 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation appendIdCriteria(new ClauseCallback() { @Override - public void onClause(Clause clause) { + public void doWithClause(Clause clause) { where.and(clause); } }, entity, id); @@ -522,134 +331,62 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation appendIdCriteria(new ClauseCallback() { @Override - public void onClause(Clause clause) { + public void doWithClause(Clause clause) { where.and(clause); } }, entity, id); } @Override - public T selectOne(Select selectQuery, Class selectClass) { - return selectOne(selectQuery.getQueryString(), selectClass); - } - - @Override - public T selectOne(String cql, Class selectClass) { - return selectOne(cql, new ReadRowCallback(cassandraConverter, selectClass)); + public T selectOne(String cql, Class type) { + return selectOne(cql, new CassandraConverterRowCallback(cassandraConverter, type)); } @Override public List update(List entities) { - String tableName = getTableName(entities.get(0).getClass()); - Assert.notNull(tableName); - return update(entities, tableName); + return update(entities, null); } @Override public List update(List entities, QueryOptions options) { - String tableName = getTableName(entities.get(0).getClass()); - Assert.notNull(tableName); - return update(entities, tableName, options); - } - - @Override - public List update(List entities, String tableName) { - return update(entities, tableName, null); - } - - @Override - public List update(List entities, String tableName, QueryOptions options) { - Assert.notNull(entities); - Assert.notEmpty(entities); - Assert.notNull(tableName); - return batchUpdate(tableName, entities, options, false); + return batchUpdate(entities, options, false); } @Override public T update(T entity) { - String tableName = getTableName(entity.getClass()); - Assert.notNull(tableName); - return update(entity, tableName); + return update(entity, null); } @Override public T update(T entity, QueryOptions options) { - String tableName = getTableName(entity.getClass()); - Assert.notNull(tableName); - return update(entity, tableName, options); - } - - @Override - public T update(T entity, String tableName) { - return update(entity, tableName, null); - } - - @Override - public T update(T entity, String tableName, QueryOptions options) { - Assert.notNull(entity); - Assert.notNull(tableName); - return update(tableName, entity, options, false); + return update(entity, options, false); } @Override public List updateAsynchronously(List entities) { - String tableName = getTableName(entities.get(0).getClass()); - Assert.notNull(tableName); - return updateAsynchronously(entities, tableName); + return updateAsynchronously(entities, null); } @Override public List updateAsynchronously(List entities, QueryOptions options) { - String tableName = getTableName(entities.get(0).getClass()); - Assert.notNull(tableName); - return updateAsynchronously(entities, tableName, options); - } - - @Override - public List updateAsynchronously(List entities, String tableName) { - return updateAsynchronously(entities, tableName, null); - } - - @Override - public List updateAsynchronously(List entities, String tableName, QueryOptions options) { - Assert.notNull(entities); - Assert.notEmpty(entities); - Assert.notNull(tableName); - return batchUpdate(tableName, entities, options, true); + return batchUpdate(entities, options, true); } @Override public T updateAsynchronously(T entity) { - String tableName = getTableName(entity.getClass()); - Assert.notNull(tableName); - return updateAsynchronously(entity, tableName); + return updateAsynchronously(entity, null); } @Override public T updateAsynchronously(T entity, QueryOptions options) { - String tableName = getTableName(entity.getClass()); - Assert.notNull(tableName); - return updateAsynchronously(entity, tableName, options); - } - - @Override - public T updateAsynchronously(T entity, String tableName) { - - return updateAsynchronously(entity, tableName, null); - } - - @Override - public T updateAsynchronously(T entity, String tableName, QueryOptions options) { - Assert.notNull(entity); - Assert.notNull(tableName); - return update(tableName, entity, options, true); + return update(entity, options, true); } /** * @param obj * @return */ - private String determineTableName(T obj) { + protected String determineTableName(T obj) { if (null != obj) { return determineTableName(obj.getClass()); } @@ -657,7 +394,7 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation return null; } - private List select(final String query, ReadRowCallback readRowCallback) { + protected List select(final String query, CassandraConverterRowCallback readRowCallback) { ResultSet resultSet = doExecute(new SessionCallback() { @@ -681,55 +418,16 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation return result; } - private Long selectCount(final Select query) { - - Long count = null; - - ResultSet resultSet = doExecute(new SessionCallback() { - - @Override - public ResultSet doInSession(Session s) throws DataAccessException { - return s.execute(query); - } - }); - - if (resultSet == null) { - return null; - } - - Iterator iterator = resultSet.iterator(); - while (iterator.hasNext()) { - Row row = iterator.next(); - count = row.getLong(0); - } - - return count; - - } - /** * @param query * @param readRowCallback * @return */ - private T selectOne(final String query, ReadRowCallback readRowCallback) { + protected T selectOne(String query, CassandraConverterRowCallback readRowCallback) { logger.info(query); - /* - * Run the Query - */ - ResultSet resultSet = doExecute(new SessionCallback() { - - @Override - public ResultSet doInSession(Session s) throws DataAccessException { - return s.execute(query); - } - }); - - if (resultSet == null) { - return null; - } + ResultSet resultSet = query(query); Iterator iterator = resultSet.iterator(); if (iterator.hasNext()) { @@ -750,63 +448,57 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation * @param tableName * @param objectToRemove */ - protected void batchDelete(final String tableName, final List entities, final QueryOptions options, - final boolean deleteAsynchronously) { + protected void batchDelete(List entities, QueryOptions options, boolean asynchronously) { Assert.notEmpty(entities); - final Batch b = createDeleteBatchQuery(tableName, entities, options, cassandraConverter); + Batch b = createDeleteBatchQuery(getTableName(entities.get(0).getClass()), entities, options, cassandraConverter); + logger.info(b.toString()); - doExecute(new SessionCallback() { + String query = b.getQueryString(); - @Override - public Object doInSession(Session s) throws DataAccessException { - - if (deleteAsynchronously) { - s.executeAsync(b); - } else { - s.execute(b); - } - - return null; - - } - }); + if (asynchronously) { + executeAsynchronously(query); + } else { + execute(query); + } } - /** - * Insert a row into a Cassandra CQL Table - * - * @param tableName - * @param entities - * @param optionsByName - * @param insertAsychronously - * @return - */ - protected List batchInsert(final String tableName, final List entities, final QueryOptions options, - final boolean insertAsychronously) { + protected T insert(T entity, QueryOptions options, boolean asynchronously) { + + Assert.notNull(entity); + + Insert insert = createInsertQuery(getTableName(entity.getClass()), entity, options, cassandraConverter); + + String query = insert.getQueryString(); + + if (asynchronously) { + executeAsynchronously(query); + } else { + execute(query); + } + + return entity; // TODO: fix this! + } + + protected List batchInsert(List entities, QueryOptions options, boolean asychronously) { Assert.notEmpty(entities); - final Batch b = createInsertBatchQuery(tableName, entities, options, cassandraConverter); - logger.info(b.getQueryString()); + Batch b = createInsertBatchQuery(getTableName(entities.get(0).getClass()), entities, options, cassandraConverter); - return doExecute(new SessionCallback>() { + String query = b.getQueryString(); + logger.info(query); - @Override - public List doInSession(Session s) throws DataAccessException { + if (asychronously) { + executeAsynchronously(query); + } else { + execute(query); + } - if (insertAsychronously) { - s.executeAsync(b); - } else { - s.execute(b); - } - - return entities; - - } - }); + return entities; // TODO: fix this! You're not supposed to necessarily return the very same entities that went in + // because the database may assign values. } /** @@ -818,93 +510,44 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation * @param updateAsychronously * @return */ - protected List batchUpdate(final String tableName, final List entities, final QueryOptions options, - final boolean updateAsychronously) { + protected List batchUpdate(List entities, QueryOptions options, boolean asychronously) { Assert.notEmpty(entities); - final Batch b = toUpdateBatchQuery(tableName, entities, options, cassandraConverter); - logger.info(b.toString()); + Batch b = toUpdateBatchQuery(getTableName(entities.get(0).getClass()), entities, options, cassandraConverter); - return doExecute(new SessionCallback>() { + String query = b.getQueryString(); + logger.info(query); - @Override - public List doInSession(Session s) throws DataAccessException { + if (asychronously) { + executeAsynchronously(query); + } else { + execute(query); + } - if (updateAsychronously) { - s.executeAsync(b); - } else { - s.execute(b); - } - - return entities; - - } - }); + return entities; // TODO: fix this! } /** * Perform the removal of a Row. * * @param tableName - * @param objectToRemove - */ - protected void delete(final String tableName, final T objectToRemove, final QueryOptions options, - final boolean deleteAsynchronously) { - - final Query q = createDeleteQuery(tableName, objectToRemove, options, cassandraConverter); - logger.info(q.toString()); - - doExecute(new SessionCallback() { - - @Override - public Object doInSession(Session s) throws DataAccessException { - - if (deleteAsynchronously) { - s.executeAsync(q); - } else { - s.execute(q); - } - - return null; - } - }); - } - - /** - * Insert a row into a Cassandra CQL Table - * - * @param tableName * @param entity */ - protected T insert(final String tableName, final T entity, final QueryOptions options, - final boolean insertAsychronously) { + protected void delete(T entity, QueryOptions options, boolean asynchronously) { - final Query q = createInsertQuery(tableName, entity, options, cassandraConverter); + Assert.notNull(entity); - logger.info(q.toString()); - if (q.getConsistencyLevel() != null) { - logger.info(q.getConsistencyLevel().name()); + Delete delete = createDeleteQuery(getTableName(entity.getClass()), entity, options, cassandraConverter); + logger.info(delete.toString()); + + String query = delete.getQueryString(); + + if (asynchronously) { + executeAsynchronously(query); + } else { + execute(query); } - if (q.getRetryPolicy() != null) { - logger.info(q.getRetryPolicy().toString()); - } - - return doExecute(new SessionCallback() { - - @Override - public T doInSession(Session s) throws DataAccessException { - - if (insertAsychronously) { - s.executeAsync(q); - } else { - s.execute(q); - } - - return entity; - - } - }); } /** @@ -916,43 +559,22 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation * @param updateAsychronously * @return */ - protected T update(final String tableName, final T entity, final QueryOptions options, - final boolean updateAsychronously) { + protected T update(T entity, QueryOptions options, boolean asychronously) { - final Query q = toUpdateQuery(tableName, entity, options, cassandraConverter); - logger.info(q.toString()); + Assert.notNull(entity); - return doExecute(new SessionCallback() { + Update q = toUpdateQuery(getTableName(entity.getClass()), entity, options, cassandraConverter); - @Override - public T doInSession(Session s) throws DataAccessException { + String query = q.getQueryString(); + logger.info(query); - if (updateAsychronously) { - s.executeAsync(q); - } else { - s.execute(q); - } - - return entity; - - } - }); - } - - /** - * Verify the object is not an iterable type - * - * @param o - */ - protected void ensureNotIterable(Object o) { - - if (null == o) { - return; + if (asychronously) { + executeAsynchronously(query); + } else { + execute(query); } - if (o.getClass().isArray() || (o instanceof Iterable) || (o instanceof Iterator)) { - throw new IllegalArgumentException("cannot use a multivalued object here."); - } + return entity; // TODO: fix this! } /** @@ -965,10 +587,10 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation * * @return The Query object to run with session.execute(); */ - public static Query createInsertQuery(String tableName, final Object objectToSave, QueryOptions options, + public static Insert createInsertQuery(String tableName, Object objectToSave, QueryOptions options, EntityWriter entityWriter) { - final Insert q = QueryBuilder.insertInto(tableName); + Insert q = QueryBuilder.insertInto(tableName); /* * Write properties @@ -988,7 +610,6 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation } return q; - } /** @@ -1001,10 +622,10 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation * * @return The Query object to run with session.execute(); */ - public static Query toUpdateQuery(String tableName, final Object objectToSave, QueryOptions options, + public static Update toUpdateQuery(String tableName, Object objectToSave, QueryOptions options, EntityWriter entityWriter) { - final Update q = QueryBuilder.update(tableName); + Update q = QueryBuilder.update(tableName); /* * Write properties @@ -1037,17 +658,17 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation * * @return The Query object to run with session.execute(); */ - public static Batch toUpdateBatchQuery(final String tableName, final List objectsToSave, QueryOptions options, + public static Batch toUpdateBatchQuery(String tableName, List objectsToSave, QueryOptions options, EntityWriter entityWriter) { /* * Return variable is a Batch statement */ - final Batch b = QueryBuilder.batch(); + Batch b = QueryBuilder.batch(); - for (final T objectToSave : objectsToSave) { + for (T objectToSave : objectsToSave) { - b.add((Statement) toUpdateQuery(tableName, objectToSave, options, entityWriter)); + b.add(toUpdateQuery(tableName, objectToSave, options, entityWriter)); } @@ -1070,13 +691,13 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation * * @return The Query object to run with session.execute(); */ - public static Batch createInsertBatchQuery(final String tableName, final List entities, QueryOptions options, + public static Batch createInsertBatchQuery(String tableName, List entities, QueryOptions options, EntityWriter entityWriter) { Batch batch = QueryBuilder.batch(); for (T entity : entities) { - batch.add((Statement) createInsertQuery(tableName, entity, options, entityWriter)); + batch.add(createInsertQuery(tableName, entity, options, entityWriter)); } CqlTemplate.addQueryOptions(batch, options); @@ -1093,7 +714,7 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation * @param optionsByName * @return */ - public static Query createDeleteQuery(String tableName, final Object object, QueryOptions options, + public static Delete createDeleteQuery(String tableName, Object object, QueryOptions options, EntityWriter entityWriter) { Delete.Selection ds = QueryBuilder.delete(); @@ -1120,22 +741,17 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation public static Batch createDeleteBatchQuery(String tableName, List entities, QueryOptions options, EntityWriter entityWriter) { + Assert.notEmpty(entities); + Assert.hasText(tableName); + Batch batch = QueryBuilder.batch(); for (T entity : entities) { - batch.add((Statement) createDeleteQuery(tableName, entity, options, entityWriter)); + batch.add(createDeleteQuery(tableName, entity, options, entityWriter)); } CqlTemplate.addQueryOptions(batch, options); return batch; } - - public boolean getUseFieldAccessOnly() { - return useFieldAccessOnly; - } - - public void setUseFieldAccessOnly(boolean useFieldAccessOnly) { - this.useFieldAccessOnly = useFieldAccessOnly; - } } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/mapping/CassandraMappingContext.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/mapping/CassandraMappingContext.java index b0e996964..126605c9a 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/mapping/CassandraMappingContext.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/mapping/CassandraMappingContext.java @@ -2,6 +2,7 @@ package org.springframework.data.cassandra.mapping; import org.springframework.cassandra.core.keyspace.CreateTableSpecification; import org.springframework.data.mapping.context.MappingContext; +import org.springframework.data.util.TypeInformation; import com.datastax.driver.core.TableMetadata; @@ -26,4 +27,24 @@ public interface CassandraMappingContext extends * @param table May not be null. */ boolean usesTable(TableMetadata table); + + /** + * Returns the {@link CassandraPersistentEntity} for the given type. If it doesn't exist, this method throws + * {@link IllegalArgumentException}. + * + * @param type The Java type of the persistent entity. + * @return The {@link CassandraPersistentEntity} describing the persistent Java type. + * @throws IllegalArgumentException if the persistent entity is unknown + */ + public CassandraPersistentEntity getRequiredPersistentEntity(Class type); + + /** + * Returns the {@link CassandraPersistentEntity} for the given type. If it doesn't exist, this method throws + * {@link IllegalArgumentException}. + * + * @param type The {@link TypeInformation} of the persistent entity. + * @return The {@link CassandraPersistentEntity} describing the persistent Java type. + * @throws IllegalArgumentException if the persistent entity is unknown + */ + public CassandraPersistentEntity getRequiredPersistentEntity(TypeInformation type); } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/mapping/DefaultCassandraMappingContext.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/mapping/DefaultCassandraMappingContext.java index c497b2541..779637227 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/mapping/DefaultCassandraMappingContext.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/mapping/DefaultCassandraMappingContext.java @@ -148,4 +148,29 @@ public class DefaultCassandraMappingContext extends return spec; } + + @Override + public CassandraPersistentEntity getRequiredPersistentEntity(Class type) { + + CassandraPersistentEntity entity = getPersistentEntity(type); + + if (entity == null) { + throw new IllegalArgumentException(String.format("no persistence metadata found for type [%s]", type.getName())); + } + + return entity; + } + + @Override + public CassandraPersistentEntity getRequiredPersistentEntity(TypeInformation type) { + + CassandraPersistentEntity entity = getPersistentEntity(type); + + if (entity == null) { + throw new IllegalArgumentException(String.format("no persistence metadata found for type [%s]", + type.getActualType())); + } + + return entity; + } } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/SimpleCassandraRepository.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/SimpleCassandraRepository.java index 3e5179475..c00a5e4f4 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/SimpleCassandraRepository.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/SimpleCassandraRepository.java @@ -19,6 +19,7 @@ import java.io.Serializable; import java.util.List; import org.springframework.cassandra.core.util.CollectionUtils; +import org.springframework.data.cassandra.core.CassandraOperations; import org.springframework.data.cassandra.core.CassandraTemplate; import org.springframework.data.cassandra.repository.CassandraRepository; import org.springframework.data.cassandra.repository.query.CassandraEntityInformation; @@ -34,7 +35,7 @@ import com.datastax.driver.core.querybuilder.Select; */ public class SimpleCassandraRepository implements CassandraRepository { - protected CassandraTemplate template; + protected CassandraOperations template; protected CassandraEntityInformation entityInformation; /** @@ -55,7 +56,7 @@ public class SimpleCassandraRepository implements Ca @Override public S save(S entity) { - return template.insert(entity, entityInformation.getTableName()); + return template.insert(entity); } @Override @@ -70,7 +71,7 @@ public class SimpleCassandraRepository implements Ca @Override public boolean exists(ID id) { - return template.countById(entityInformation.getJavaType(), id) >= 1; // TODO: == instead of >= ? + return template.exists(entityInformation.getJavaType(), id); } @Override @@ -109,6 +110,6 @@ public class SimpleCassandraRepository implements Ca } protected List findAll(Select query) { - return template.select(query, entityInformation.getJavaType()); + return template.select(query.getQueryString(), entityInformation.getJavaType()); } } diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/test/integration/template/CassandraDataOperationsTest.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/test/integration/template/CassandraDataOperationsTest.java index 56ae82be3..b32298259 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/test/integration/template/CassandraDataOperationsTest.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/test/integration/template/CassandraDataOperationsTest.java @@ -61,7 +61,7 @@ import com.datastax.driver.core.querybuilder.Select; public class CassandraDataOperationsTest { @Autowired - private CassandraOperations cassandraTemplate; + private CassandraOperations template; private static Logger log = LoggerFactory.getLogger(CassandraDataOperationsTest.class); @@ -92,16 +92,13 @@ public class CassandraDataOperationsTest { @Test public void insertTest() { - /* - * Test Single Insert with entity - */ Book b1 = new Book(); b1.setIsbn("123456-1"); b1.setTitle("Spring Data Cassandra Guide"); b1.setAuthor("Cassandra Guru"); b1.setPages(521); - cassandraTemplate.insert(b1); + template.insert(b1); Book b2 = new Book(); b2.setIsbn("123456-2"); @@ -109,11 +106,8 @@ public class CassandraDataOperationsTest { b2.setAuthor("Cassandra Guru"); b2.setPages(521); - cassandraTemplate.insert(b2, "book_alt"); + template.insert(b2); - /* - * Test Single Insert with entity - */ Book b3 = new Book(); b3.setIsbn("123456-3"); b3.setTitle("Spring Data Cassandra Guide"); @@ -124,34 +118,27 @@ public class CassandraDataOperationsTest { options.setConsistencyLevel(ConsistencyLevel.ONE); options.setRetryPolicy(RetryPolicy.DOWNGRADING_CONSISTENCY); - cassandraTemplate.insert(b3, "book", options); + template.insert(b3, options); - /* - * Test Single Insert with entity - */ Book b5 = new Book(); b5.setIsbn("123456-5"); b5.setTitle("Spring Data Cassandra Guide"); b5.setAuthor("Cassandra Guru"); b5.setPages(265); - cassandraTemplate.insert(b5, options); - + template.insert(b5, options); } @Test public void insertAsynchronouslyTest() { - /* - * Test Single Insert with entity - */ Book b1 = new Book(); b1.setIsbn("123456-1"); b1.setTitle("Spring Data Cassandra Guide"); b1.setAuthor("Cassandra Guru"); b1.setPages(521); - cassandraTemplate.insertAsynchronously(b1); + template.insertAsynchronously(b1); Book b2 = new Book(); b2.setIsbn("123456-2"); @@ -159,7 +146,7 @@ public class CassandraDataOperationsTest { b2.setAuthor("Cassandra Guru"); b2.setPages(521); - cassandraTemplate.insertAsynchronously(b2, "book_alt"); + template.insertAsynchronously(b2); /* * Test Single Insert with entity @@ -174,7 +161,7 @@ public class CassandraDataOperationsTest { options.setConsistencyLevel(ConsistencyLevel.ONE); options.setRetryPolicy(RetryPolicy.DOWNGRADING_CONSISTENCY); - cassandraTemplate.insertAsynchronously(b3, "book", options); + template.insertAsynchronously(b3, options); /* * Test Single Insert with entity @@ -194,7 +181,7 @@ public class CassandraDataOperationsTest { b5.setAuthor("Cassandra Guru"); b5.setPages(265); - cassandraTemplate.insertAsynchronously(b5, options); + template.insertAsynchronously(b5, options); } @@ -209,19 +196,19 @@ public class CassandraDataOperationsTest { books = getBookList(20); - cassandraTemplate.insert(books); + template.insert(books); books = getBookList(20); - cassandraTemplate.insert(books, "book_alt"); + template.insert(books); books = getBookList(20); - cassandraTemplate.insert(books, "book", options); + template.insert(books, options); books = getBookList(20); - cassandraTemplate.insert(books, options); + template.insert(books, options); } @@ -236,19 +223,19 @@ public class CassandraDataOperationsTest { books = getBookList(20); - cassandraTemplate.insertAsynchronously(books); + template.insertAsynchronously(books); books = getBookList(20); - cassandraTemplate.insertAsynchronously(books, "book_alt"); + template.insertAsynchronously(books); books = getBookList(20); - cassandraTemplate.insertAsynchronously(books, "book", options); + template.insertAsynchronously(books, options); books = getBookList(20); - cassandraTemplate.insertAsynchronously(books, options); + template.insertAsynchronously(books, options); } @@ -290,7 +277,7 @@ public class CassandraDataOperationsTest { b1.setAuthor("Cassandra Guru"); b1.setPages(521); - cassandraTemplate.update(b1); + template.update(b1); Book b2 = new Book(); b2.setIsbn("123456-2"); @@ -298,7 +285,7 @@ public class CassandraDataOperationsTest { b2.setAuthor("Cassandra Guru"); b2.setPages(521); - cassandraTemplate.update(b2, "book_alt"); + template.update(b2); /* * Test Single Insert with entity @@ -309,7 +296,7 @@ public class CassandraDataOperationsTest { b3.setAuthor("Cassandra Guru"); b3.setPages(265); - cassandraTemplate.update(b3, "book", options); + template.update(b3, options); /* * Test Single Insert with entity @@ -320,7 +307,7 @@ public class CassandraDataOperationsTest { b5.setAuthor("Cassandra Guru"); b5.setPages(265); - cassandraTemplate.update(b5, options); + template.update(b5, options); } @@ -342,7 +329,7 @@ public class CassandraDataOperationsTest { b1.setAuthor("Cassandra Guru"); b1.setPages(521); - cassandraTemplate.updateAsynchronously(b1); + template.updateAsynchronously(b1); Book b2 = new Book(); b2.setIsbn("123456-2"); @@ -350,7 +337,7 @@ public class CassandraDataOperationsTest { b2.setAuthor("Cassandra Guru"); b2.setPages(521); - cassandraTemplate.updateAsynchronously(b2, "book_alt"); + template.updateAsynchronously(b2); /* * Test Single Insert with entity @@ -361,7 +348,7 @@ public class CassandraDataOperationsTest { b3.setAuthor("Cassandra Guru"); b3.setPages(265); - cassandraTemplate.updateAsynchronously(b3, "book", options); + template.updateAsynchronously(b3, options); /* * Test Single Insert with entity @@ -372,7 +359,7 @@ public class CassandraDataOperationsTest { b5.setAuthor("Cassandra Guru"); b5.setPages(265); - cassandraTemplate.updateAsynchronously(b5, options); + template.updateAsynchronously(b5, options); } @@ -387,35 +374,35 @@ public class CassandraDataOperationsTest { books = getBookList(20); - cassandraTemplate.insert(books); + template.insert(books); alterBooks(books); - cassandraTemplate.update(books); + template.update(books); books = getBookList(20); - cassandraTemplate.insert(books, "book_alt"); + template.insert(books); alterBooks(books); - cassandraTemplate.update(books, "book_alt"); + template.update(books); books = getBookList(20); - cassandraTemplate.insert(books, "book", options); + template.insert(books, options); alterBooks(books); - cassandraTemplate.update(books, "book", options); + template.update(books, options); books = getBookList(20); - cassandraTemplate.insert(books, options); + template.insert(books, options); alterBooks(books); - cassandraTemplate.update(books, options); + template.update(books, options); } @@ -430,35 +417,35 @@ public class CassandraDataOperationsTest { books = getBookList(20); - cassandraTemplate.insert(books); + template.insert(books); alterBooks(books); - cassandraTemplate.updateAsynchronously(books); + template.updateAsynchronously(books); books = getBookList(20); - cassandraTemplate.insert(books, "book_alt"); + template.insert(books); alterBooks(books); - cassandraTemplate.updateAsynchronously(books, "book_alt"); + template.updateAsynchronously(books); books = getBookList(20); - cassandraTemplate.insert(books, "book", options); + template.insert(books, options); alterBooks(books); - cassandraTemplate.updateAsynchronously(books, "book", options); + template.updateAsynchronously(books, options); books = getBookList(20); - cassandraTemplate.insert(books, options); + template.insert(books, options); alterBooks(books); - cassandraTemplate.updateAsynchronously(books, options); + template.updateAsynchronously(books, options); } @@ -489,12 +476,12 @@ public class CassandraDataOperationsTest { Book b1 = new Book(); b1.setIsbn("123456-1"); - cassandraTemplate.delete(b1); + template.delete(b1); Book b2 = new Book(); b2.setIsbn("123456-2"); - cassandraTemplate.delete(b2, "book_alt"); + template.delete(b2); /* * Test Single Insert with entity @@ -502,7 +489,7 @@ public class CassandraDataOperationsTest { Book b3 = new Book(); b3.setIsbn("123456-3"); - cassandraTemplate.delete(b3, "book", options); + template.delete(b3, options); /* * Test Single Insert with entity @@ -510,7 +497,7 @@ public class CassandraDataOperationsTest { Book b5 = new Book(); b5.setIsbn("123456-5"); - cassandraTemplate.delete(b5, options); + template.delete(b5, options); } @@ -529,12 +516,12 @@ public class CassandraDataOperationsTest { Book b1 = new Book(); b1.setIsbn("123456-1"); - cassandraTemplate.deleteAsynchronously(b1); + template.deleteAsynchronously(b1); Book b2 = new Book(); b2.setIsbn("123456-2"); - cassandraTemplate.deleteAsynchronously(b2, "book_alt"); + template.deleteAsynchronously(b2); /* * Test Single Insert with entity @@ -542,7 +529,7 @@ public class CassandraDataOperationsTest { Book b3 = new Book(); b3.setIsbn("123456-3"); - cassandraTemplate.deleteAsynchronously(b3, "book", options); + template.deleteAsynchronously(b3, options); /* * Test Single Insert with entity @@ -550,7 +537,7 @@ public class CassandraDataOperationsTest { Book b5 = new Book(); b5.setIsbn("123456-5"); - cassandraTemplate.deleteAsynchronously(b5, options); + template.deleteAsynchronously(b5, options); } @@ -565,27 +552,27 @@ public class CassandraDataOperationsTest { books = getBookList(20); - cassandraTemplate.insert(books); + template.insert(books); - cassandraTemplate.delete(books); + template.delete(books); books = getBookList(20); - cassandraTemplate.insert(books, "book_alt"); + template.insert(books); - cassandraTemplate.delete(books, "book_alt"); + template.delete(books); books = getBookList(20); - cassandraTemplate.insert(books, "book", options); + template.insert(books, options); - cassandraTemplate.delete(books, "book", options); + template.delete(books, options); books = getBookList(20); - cassandraTemplate.insert(books, options); + template.insert(books, options); - cassandraTemplate.delete(books, options); + template.delete(books, options); } @@ -600,27 +587,27 @@ public class CassandraDataOperationsTest { books = getBookList(20); - cassandraTemplate.insert(books); + template.insert(books); - cassandraTemplate.deleteAsynchronously(books); + template.deleteAsynchronously(books); books = getBookList(20); - cassandraTemplate.insert(books, "book_alt"); + template.insert(books); - cassandraTemplate.deleteAsynchronously(books, "book_alt"); + template.deleteAsynchronously(books); books = getBookList(20); - cassandraTemplate.insert(books, "book", options); + template.insert(books, options); - cassandraTemplate.deleteAsynchronously(books, "book", options); + template.deleteAsynchronously(books, options); books = getBookList(20); - cassandraTemplate.insert(books, options); + template.insert(books, options); - cassandraTemplate.deleteAsynchronously(books, options); + template.deleteAsynchronously(books, options); } @@ -636,12 +623,12 @@ public class CassandraDataOperationsTest { b1.setAuthor("Cassandra Guru"); b1.setPages(521); - cassandraTemplate.insert(b1); + template.insert(b1); Select select = QueryBuilder.select().all().from("book"); select.where(QueryBuilder.eq("isbn", "123456-1")); - Book b = cassandraTemplate.selectOne(select, Book.class); + Book b = template.selectOne(select.getQueryString(), Book.class); log.info("SingleSelect Book Title -> " + b.getTitle()); log.info("SingleSelect Book Author -> " + b.getAuthor()); @@ -656,11 +643,11 @@ public class CassandraDataOperationsTest { List books = getBookList(20); - cassandraTemplate.insert(books); + template.insert(books); Select select = QueryBuilder.select().all().from("book"); - List b = cassandraTemplate.select(select, Book.class); + List b = template.select(select.getQueryString(), Book.class); log.info("Book Count -> " + b.size()); @@ -671,23 +658,16 @@ public class CassandraDataOperationsTest { @Test public void selectCountTest() { - List books = getBookList(20); + int count = 20; + List books = getBookList(count); - cassandraTemplate.insert(books); - - Select select = QueryBuilder.select().countAll().from("book"); - - Long count = cassandraTemplate.count(select); - - log.info("Book Count -> " + count); - - Assert.assertEquals(count, new Long(20)); + template.insert(books); + Assert.assertEquals(count, template.count(Book.class)); } @After public void clearCassandra() { EmbeddedCassandraServerHelper.cleanEmbeddedCassandra(); - } }