From 7250d0cc78d67cd918164df0e8b56c64779b337e Mon Sep 17 00:00:00 2001 From: dwebb Date: Wed, 13 Nov 2013 09:30:47 -0500 Subject: [PATCH] IN PROGRESS - issue DATACASS-32: Implement the TemplateAPI for CQL https://jira.springsource.org/browse/DATACASS-32 Added executeQueryAsyn to Operations and Template. Modified executeQuery(String) to use the SessionCallback. --- .../cassandra/core/CassandraOperations.java | 92 +++-- .../cassandra/core/CassandraTemplate.java | 377 ++++++++++-------- 2 files changed, 256 insertions(+), 213 deletions(-) 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 da8fe108e..6200e107e 100644 --- a/src/main/java/org/springframework/data/cassandra/core/CassandraOperations.java +++ b/src/main/java/org/springframework/data/cassandra/core/CassandraOperations.java @@ -21,20 +21,20 @@ import org.springframework.data.cassandra.convert.CassandraConverter; import org.springframework.data.cassandra.dto.RingMember; import com.datastax.driver.core.ResultSet; +import com.datastax.driver.core.ResultSetFuture; /** * @author Alex Shvid */ public interface CassandraOperations { - /** * Describe the current Ring * * @return The list of ring tokens that are active in the cluster */ List describeRing(); - + /** * The table name used for the specified class by this template. * @@ -42,15 +42,23 @@ public interface CassandraOperations { * @return */ String getTableName(Class entityClass); - + /** * Execute query and return Cassandra ResultSet * * @param query must not be {@literal null}. * @return */ - ResultSet executeQuery(String query); - + ResultSet executeQuery(final String query); + + /** + * Execute async query and return Cassandra ResultSetFuture + * + * @param query must not be {@literal null}. + * @return + */ + ResultSetFuture executeQueryAsync(final String query); + /** * Execute query and convert ResultSet to the list of entities * @@ -58,8 +66,8 @@ public interface CassandraOperations { * @param selectClass must not be {@literal null}, mapped entity type. * @return */ - List select(String query, Class selectClass); - + List select(String query, Class selectClass); + /** * Execute query and convert ResultSet to the entity * @@ -67,22 +75,22 @@ public interface CassandraOperations { * @param selectClass must not be {@literal null}, mapped entity type. * @return */ - T selectOne(String query, Class selectClass); - - /** - * Insert the given object to the table by id. - * - * @param object - */ - void insert(Object entity); + T selectOne(String query, Class selectClass); /** * Insert the given object to the table by id. * * @param object */ - void insert(Object entity, String tableName); - + void insert(Object entity); + + /** + * Insert the given object to the table by id. + * + * @param object + */ + void insert(Object entity, String tableName); + /** * Remove the given object from the table by id. * @@ -97,55 +105,55 @@ public interface CassandraOperations { * @param table must not be {@literal null} or empty. */ void remove(Object object, String tableName); - + /** * Create a table with the name and fields indicated by the entity class * * @param entityClass class that determines metadata of the table to create/drop. */ - void createTable(Class entityClass); - + void createTable(Class entityClass); + /** * Create a table with the name and fields indicated by the entity class * * @param entityClass class that determines metadata of the table to create/drop. * @param tableName explicit name of the table */ - void createTable(Class entityClass, String tableName); - - /** - * Alter table with the name and fields indicated by the entity class - * - * @param entityClass class that determines metadata of the table to create/drop. - */ - void alterTable(Class entityClass); - - /** - * Alter table with the name and fields indicated by the entity class - * - * @param entityClass class that determines metadata of the table to create/drop. - * @param tableName explicit name of the table - */ - void alterTable(Class entityClass, String tableName); + void createTable(Class entityClass, String tableName); /** * Alter table with the name and fields indicated by the entity class * * @param entityClass class that determines metadata of the table to create/drop. - */ - void dropTable(Class entityClass); + */ + void alterTable(Class entityClass); + + /** + * Alter table with the name and fields indicated by the entity class + * + * @param entityClass class that determines metadata of the table to create/drop. + * @param tableName explicit name of the table + */ + void alterTable(Class entityClass, String tableName); + + /** + * Alter table with the name and fields indicated by the entity class + * + * @param entityClass class that determines metadata of the table to create/drop. + */ + void dropTable(Class entityClass); /** * Alter table with the name and fields indicated by the entity class * * @param tableName explicit name of the table. - */ - void dropTable(String tableName); - + */ + void dropTable(String tableName); + /** * Returns the underlying {@link CassandraConverter}. * * @return */ - CassandraConverter getConverter(); + CassandraConverter getConverter(); } 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 15cdf787f..d85b657ee 100644 --- a/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java +++ b/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java @@ -44,6 +44,7 @@ import com.datastax.driver.core.Host; import com.datastax.driver.core.Metadata; import com.datastax.driver.core.Query; import com.datastax.driver.core.ResultSet; +import com.datastax.driver.core.ResultSetFuture; import com.datastax.driver.core.Row; import com.datastax.driver.core.Session; import com.datastax.driver.core.exceptions.NoHostAvailableException; @@ -52,27 +53,27 @@ import com.datastax.driver.core.exceptions.NoHostAvailableException; * @author Alex Shvid */ public class CassandraTemplate implements CassandraOperations { - + private static Logger log = LoggerFactory.getLogger(CassandraTemplate.class); - - private static final Collection ITERABLE_CLASSES; - static { - Set iterableClasses = new HashSet(); - iterableClasses.add(List.class.getName()); - iterableClasses.add(Collection.class.getName()); - iterableClasses.add(Iterator.class.getName()); + private static final Collection ITERABLE_CLASSES; + static { - ITERABLE_CLASSES = Collections.unmodifiableCollection(iterableClasses); - } + Set iterableClasses = new HashSet(); + iterableClasses.add(List.class.getName()); + iterableClasses.add(Collection.class.getName()); + iterableClasses.add(Iterator.class.getName()); - private final Keyspace keyspace; + 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; /** @@ -86,14 +87,14 @@ public class CassandraTemplate implements CassandraOperations { this.cassandraConverter = keyspace.getCassandraConverter(); this.mappingContext = this.cassandraConverter.getMappingContext(); } - - /** - * @param classLoader - */ - public void setBeanClassLoader(ClassLoader classLoader) { - this.beanClassLoader = classLoader; + + /** + * @param classLoader + */ + public void setBeanClassLoader(ClassLoader classLoader) { + this.beanClassLoader = classLoader; } - + /* (non-Javadoc) * @see org.springframework.data.cassandra.core.CassandraOperations#describeRing() */ @@ -104,7 +105,7 @@ public class CassandraTemplate implements CassandraOperations { * Initialize the return variable */ List ring = new ArrayList(); - + /* * Get the cluster metadata for this session */ @@ -114,33 +115,33 @@ public class CassandraTemplate implements CassandraOperations { public Metadata doInSession(Session s) throws DataAccessException { return s.getCluster().getMetadata(); } - + }); - + /* * Get all hosts in the cluster */ Set hosts = clusterMetadata.getAllHosts(); - + /* * Loop variables */ RingMember member = null; - + /* * Populate Ring with Host Metadata */ - for (Host h: hosts) { - + for (Host h : hosts) { + member = new RingMember(h); ring.add(member); } - + /* * Return */ return ring; - + } /** @@ -150,7 +151,7 @@ public class CassandraTemplate implements CassandraOperations { * @return */ protected CassandraPersistentEntity getEntity(Object o) { - + CassandraPersistentEntity entity = null; try { String entityClassName = o.getClass().getName(); @@ -160,29 +161,55 @@ public class CassandraTemplate implements CassandraOperations { e.printStackTrace(); } catch (LinkageError e) { e.printStackTrace(); - } finally {} - + } finally { + } + return entity; - + } + /* (non-Javadoc) * @see org.springframework.data.cassandra.core.CassandraOperations#getTableName(java.lang.Class) */ public String getTableName(Class entityClass) { return determineTableName(entityClass); } - + /* (non-Javadoc) * @see org.springframework.data.cassandra.core.CassandraOperations#executeQuery(java.lang.String) */ - public ResultSet executeQuery(String query) { - try { - return session.execute(query); - } catch (NoHostAvailableException e) { - throw new CassandraConnectionFailureException("no host available", e); - } catch (RuntimeException e) { - throw potentiallyConvertRuntimeException(e); - } + public ResultSet executeQuery(final String query) { + + return execute(new SessionCallback() { + + @Override + public ResultSet doInSession(Session s) throws DataAccessException { + + return s.execute(query); + + } + + }); + + } + + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#executeQueryAsync(java.lang.String) + */ + @Override + public ResultSetFuture executeQueryAsync(final String query) { + + return execute(new SessionCallback() { + + @Override + public ResultSetFuture doInSession(Session s) throws DataAccessException { + + return s.executeAsync(query); + + } + + }); + } /* (non-Javadoc) @@ -191,21 +218,21 @@ public class CassandraTemplate implements CassandraOperations { public List select(String query, Class selectClass) { return selectInternal(query, new ReadRowCallback(cassandraConverter, selectClass)); } - + /* (non-Javadoc) * @see org.springframework.data.cassandra.core.CassandraOperations#selectOne(java.lang.String, java.lang.Class) */ public T selectOne(String query, Class selectClass) { return selectOneInternal(query, new ReadRowCallback(cassandraConverter, selectClass)); } - + /* (non-Javadoc) * @see org.springframework.data.cassandra.core.CassandraOperations#getConverter() */ public CassandraConverter getConverter() { return cassandraConverter; } - + /** * Simple {@link RowCallback} that will transform {@link Row} into the given target type using the given * {@link EntityReader}. @@ -228,8 +255,8 @@ public class CassandraTemplate implements CassandraOperations { T source = reader.read(type, object); return source; } - } - + } + /** * @param query * @param readRowCallback @@ -240,7 +267,7 @@ public class CassandraTemplate implements CassandraOperations { ResultSet resultSet = session.execute(query); List result = new ArrayList(); Iterator iterator = resultSet.iterator(); - while(iterator.hasNext()) { + while (iterator.hasNext()) { Row row = iterator.next(); result.add(readRowCallback.doWith(row)); } @@ -251,7 +278,7 @@ public class CassandraTemplate implements CassandraOperations { throw potentiallyConvertRuntimeException(e); } } - + T selectOneInternal(String query, ReadRowCallback readRowCallback) { try { ResultSet resultSet = session.execute(query); @@ -271,19 +298,19 @@ public class CassandraTemplate implements CassandraOperations { throw potentiallyConvertRuntimeException(e); } } - - /** - * @param obj - * @return - */ - private String determineTableName(T obj) { - if (null != obj) { - return determineTableName(obj.getClass()); - } - return null; -} - + /** + * @param obj + * @return + */ + private String determineTableName(T obj) { + if (null != obj) { + return determineTableName(obj.getClass()); + } + + return null; + } + /** * @param entityClass * @return @@ -302,47 +329,47 @@ public class CassandraTemplate implements CassandraOperations { } return entity.getTable(); } - - private RuntimeException potentiallyConvertRuntimeException( - RuntimeException ex) { + + 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 objectToSave - */ - protected T doInsert(final String tableName, final T objectToSave) { + /** + * Insert a row into a Cassandra CQL Table + * + * @param tableName + * @param objectToSave + */ + protected T doInsert(final String tableName, final T objectToSave) { + + CassandraPersistentEntity entity = getEntity(objectToSave); + + Assert.notNull(entity); + + try { - CassandraPersistentEntity entity = getEntity(objectToSave); - - Assert.notNull(entity); - - try { - final Query q = CQLUtils.toInsertQuery(keyspace.getKeyspace(), tableName, objectToSave, entity); log.info(q.toString()); - - return execute(new SessionCallback() { - - public T doInSession(Session s) throws DataAccessException { - - s.execute(q); - + + return execute(new SessionCallback() { + + public T doInSession(Session s) throws DataAccessException { + + s.execute(q); + return objectToSave; - + } }); - - } catch (EntityWriterException e) { - throw exceptionTranslator.translateExceptionIfPossible(new RuntimeException("Failed to translate Object to Query", e)); - } - - } - + + } catch (EntityWriterException e) { + throw exceptionTranslator.translateExceptionIfPossible(new RuntimeException( + "Failed to translate Object to Query", e)); + } + + } + /** * Execute a command at the Session Level * @@ -350,55 +377,55 @@ public class CassandraTemplate implements CassandraOperations { * @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) - */ - public void insert(Object objectToSave) { - ensureNotIterable(objectToSave); - insert(objectToSave, determineTableName(objectToSave)); - } + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#insert(java.lang.Object) + */ + public void insert(Object objectToSave) { + ensureNotIterable(objectToSave); + insert(objectToSave, determineTableName(objectToSave)); + } - /* (non-Javadoc) - * @see org.springframework.data.cassandra.core.CassandraOperations#insert(java.lang.Object, java.lang.String) - */ - public void insert(Object objectToSave, String tableName) { - ensureNotIterable(objectToSave); - doInsert(tableName, objectToSave); - } - - /** - * 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."); - } - } - } + /* (non-Javadoc) + * @see org.springframework.data.cassandra.core.CassandraOperations#insert(java.lang.Object, java.lang.String) + */ + public void insert(Object objectToSave, String tableName) { + ensureNotIterable(objectToSave); + doInsert(tableName, objectToSave); + } + /** + * 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."); + } + } + } /* (non-Javadoc) * @see org.springframework.data.cassandra.core.CassandraOperations#remove(java.lang.Object) */ @Override public void remove(Object object) { - + remove(object, determineTableName(object.getClass())); - + } /* (non-Javadoc) @@ -408,36 +435,43 @@ public class CassandraTemplate implements CassandraOperations { public void remove(Object object, String tableName) { CassandraPersistentEntity entityClass = getEntity(object); - + Assert.notNull(entityClass); - + doRemove(object, tableName); - + } - + + /** + * Perform the removal of a Row. + * + * @param objectToRemove + * @param tableName + */ protected void doRemove(final Object objectToRemove, final String tableName) { - - CassandraPersistentEntity entity = getEntity(objectToRemove); - - Assert.notNull(entity); - - try { - + + 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 ResultSet doInSession(Session s) throws DataAccessException { - - return s.execute(q); - + + execute(new SessionCallback() { + + public ResultSet doInSession(Session s) throws DataAccessException { + + return s.execute(q); + } }); - - } catch (EntityWriterException e) { - throw exceptionTranslator.translateExceptionIfPossible(new RuntimeException("Failed to translate Object to Query", e)); - } + + } catch (EntityWriterException e) { + throw exceptionTranslator.translateExceptionIfPossible(new RuntimeException( + "Failed to translate Object to Query", e)); + } } /* (non-Javadoc) @@ -446,18 +480,18 @@ public class CassandraTemplate implements CassandraOperations { @Override public void createTable(Class entityClass) { - try { - + try { + final CassandraPersistentEntity entity = mappingContext.getPersistentEntity(entityClass); final String useTableName = entity.getTable(); - - createTable(entityClass, useTableName); + + createTable(entityClass, useTableName); } catch (LinkageError e) { e.printStackTrace(); - } finally {} + } finally { + } - } /* (non-Javadoc) @@ -466,29 +500,30 @@ public class CassandraTemplate implements CassandraOperations { @Override public void createTable(Class entityClass, final String tableName) { - try { - + try { + final CassandraPersistentEntity entity = mappingContext.getPersistentEntity(entityClass); - - execute(new SessionCallback() { - - public Object doInSession(Session s) throws DataAccessException { - - String cql = CQLUtils.createTable(tableName, entity); - - log.info("CREATE TABLE CQL -> " + cql); - + + execute(new SessionCallback() { + + public Object doInSession(Session s) throws DataAccessException { + + String cql = CQLUtils.createTable(tableName, entity); + + log.info("CREATE TABLE CQL -> " + cql); + s.execute(cql); - + return null; - + } }); - + } catch (LinkageError e) { e.printStackTrace(); - } finally {} - + } finally { + } + } /* (non-Javadoc) @@ -497,7 +532,7 @@ public class CassandraTemplate implements CassandraOperations { @Override public void alterTable(Class entityClass) { // TODO Auto-generated method stub - + } /* (non-Javadoc) @@ -506,7 +541,7 @@ public class CassandraTemplate implements CassandraOperations { @Override public void alterTable(Class entityClass, String tableName) { // TODO Auto-generated method stub - + } /* (non-Javadoc) @@ -515,7 +550,7 @@ public class CassandraTemplate implements CassandraOperations { @Override public void dropTable(Class entityClass) { // TODO Auto-generated method stub - + } /* (non-Javadoc) @@ -524,7 +559,7 @@ public class CassandraTemplate implements CassandraOperations { @Override public void dropTable(String tableName) { // TODO Auto-generated method stub - + } }