From 810b217fd761318f39f2391e1873bd84e4dd74c7 Mon Sep 17 00:00:00 2001 From: dwebb Date: Wed, 13 Nov 2013 16:50:43 -0500 Subject: [PATCH] wip: Completed implementation of inserts and deletes. --- .../cassandra/core/CassandraTemplate.java | 186 +++++++++++++++--- .../data/cassandra/util/CqlUtils.java | 66 +++++++ 2 files changed, 229 insertions(+), 23 deletions(-) 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 b1eb2fb2f..ec0cad58d 100644 --- a/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java +++ b/src/main/java/org/springframework/data/cassandra/core/CassandraTemplate.java @@ -47,6 +47,7 @@ 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; +import com.datastax.driver.core.querybuilder.Batch; /** * The Cassandra Template is a convenience API for all Cassnadta DML Operations. @@ -346,7 +347,47 @@ public class CassandraTemplate implements CassandraOperations { * @param tableName * @param entity */ - protected T doInsert(final String tableName, final T entity) { + protected List doBatchInsert(final String tableName, final List entities, final boolean insertAsychronously) { + + Assert.notEmpty(entities); + + CassandraPersistentEntity CPEntity = getEntity(entities.get(0)); + + Assert.notNull(CPEntity); + + try { + + final Batch b = CqlUtils.toInsertBatchQuery(keyspace.getKeyspace(), tableName, entities, CPEntity); + log.info(b.toString()); + + return execute(new SessionCallback>() { + + public List doInSession(Session s) throws DataAccessException { + + if (insertAsychronously) { + s.executeAsync(b); + } else { + s.execute(b); + } + + return entities; + + } + }); + + } catch (EntityWriterException e) { + throw exceptionTranslator.translateExceptionIfPossible(new RuntimeException( + "Failed to translate Object to Query", e)); + } + } + + /** + * Insert a row into a Cassandra CQL Table + * + * @param tableName + * @param entity + */ + protected T doInsert(final String tableName, final T entity, final boolean insertAsychronously) { CassandraPersistentEntity CPEntity = getEntity(entity); @@ -361,7 +402,11 @@ public class CassandraTemplate implements CassandraOperations { public T doInSession(Session s) throws DataAccessException { - s.execute(q); + if (insertAsychronously) { + s.executeAsync(q); + } else { + s.execute(q); + } return entity; @@ -394,7 +439,7 @@ public class CassandraTemplate implements CassandraOperations { * @param objectToRemove * @param tableName */ - protected void doRemove(final Object objectToRemove, final String tableName) { + protected void doDelete(final Object objectToRemove, final String tableName, final boolean deleteAsynchronously) { CassandraPersistentEntity entity = getEntity(objectToRemove); @@ -405,11 +450,57 @@ public class CassandraTemplate implements CassandraOperations { final Query q = CqlUtils.toDeleteQuery(keyspace.getKeyspace(), tableName, objectToRemove, entity); log.info(q.toString()); - execute(new SessionCallback() { + execute(new SessionCallback() { - public ResultSet doInSession(Session s) throws DataAccessException { + public Object doInSession(Session s) throws DataAccessException { - return s.execute(q); + if (deleteAsynchronously) { + s.executeAsync(q); + } else { + s.execute(q); + } + + return null; + + } + }); + + } catch (EntityWriterException e) { + throw exceptionTranslator.translateExceptionIfPossible(new RuntimeException( + "Failed to translate Object to Query", e)); + } + } + + /** + * Perform the deletion on a list of objects + * + * @param objectToRemove + * @param tableName + */ + protected void doBatchDelete(final String tableName, final List entities, final boolean deleteAsynchronously) { + + Assert.notEmpty(entities); + + CassandraPersistentEntity CPEntity = getEntity(entities.get(0)); + + Assert.notNull(CPEntity); + + try { + + final Batch b = CqlUtils.toDeleteBatchQuery(keyspace.getKeyspace(), tableName, entities, CPEntity); + log.info(b.toString()); + + execute(new SessionCallback() { + + public Object doInSession(Session s) throws DataAccessException { + + if (deleteAsynchronously) { + s.executeAsync(b); + } else { + s.execute(b); + } + + return null; } }); @@ -445,7 +536,12 @@ public class CassandraTemplate implements CassandraOperations { @Override public T insert(T entity) { ensureNotIterable(entity); - return insert(entity, determineTableName(entity)); + + String tableName = determineTableName(entity); + + Assert.notNull(tableName); + + return insert(entity, tableName); } /* (non-Javadoc) @@ -453,8 +549,15 @@ public class CassandraTemplate implements CassandraOperations { */ @Override public List insert(List entities) { - // TODO Auto-generated method stub - return null; + + Assert.notNull(entities); + Assert.notEmpty(entities); + + String tableName = getTableName(entities.get(0).getClass()); + + Assert.notNull(tableName); + + return insert(entities, tableName); } /* (non-Javadoc) @@ -463,7 +566,7 @@ public class CassandraTemplate implements CassandraOperations { @Override public T insert(T entity, String tableName) { ensureNotIterable(entity); - return doInsert(tableName, entity); + return doInsert(tableName, entity, false); } /* (non-Javadoc) @@ -471,8 +574,13 @@ public class CassandraTemplate implements CassandraOperations { */ @Override public List insert(List entities, String tableName) { - // TODO Auto-generated method stub - return null; + + Assert.notNull(entities); + Assert.notEmpty(entities); + Assert.notNull(tableName); + + return doBatchInsert(tableName, entities, false); + } /* (non-Javadoc) @@ -480,8 +588,14 @@ public class CassandraTemplate implements CassandraOperations { */ @Override public T insertAsynchronously(T entity) { - // TODO Auto-generated method stub - return null; + + ensureNotIterable(entity); + + String tableName = determineTableName(entity); + + Assert.notNull(tableName); + + return insertAsynchronously(entity, tableName); } /* (non-Javadoc) @@ -489,8 +603,15 @@ public class CassandraTemplate implements CassandraOperations { */ @Override public List insertAsynchronously(List entities) { - // TODO Auto-generated method stub - return null; + + Assert.notNull(entities); + Assert.notEmpty(entities); + + String tableName = getTableName(entities.get(0).getClass()); + + Assert.notNull(tableName); + + return insertAsynchronously(entities, tableName); } /* (non-Javadoc) @@ -498,8 +619,10 @@ public class CassandraTemplate implements CassandraOperations { */ @Override public T insertAsynchronously(T entity, String tableName) { - // TODO Auto-generated method stub - return null; + + ensureNotIterable(entity); + + return doInsert(tableName, entity, true); } /* (non-Javadoc) @@ -507,8 +630,13 @@ public class CassandraTemplate implements CassandraOperations { */ @Override public List insertAsynchronously(List entities, String tableName) { - // TODO Auto-generated method stub - return null; + + Assert.notNull(entities); + Assert.notEmpty(entities); + Assert.notNull(tableName); + + return doBatchInsert(tableName, entities, true); + } /* (non-Javadoc) @@ -524,7 +652,15 @@ public class CassandraTemplate implements CassandraOperations { */ @Override public void delete(List entities) { - // TODO Auto-generated method stub + + Assert.notNull(entities); + Assert.notEmpty(entities); + + String tableName = getTableName(entities.get(0).getClass()); + + Assert.notNull(tableName); + + delete(entities, tableName); } @@ -538,7 +674,7 @@ public class CassandraTemplate implements CassandraOperations { Assert.notNull(entityClass); - doRemove(entity, tableName); + doDelete(entity, tableName, false); } @@ -547,8 +683,12 @@ public class CassandraTemplate implements CassandraOperations { */ @Override public void delete(List entities, String tableName) { - // TODO Auto-generated method stub + Assert.notNull(entities); + Assert.notEmpty(entities); + Assert.notNull(tableName); + + doBatchDelete(tableName, entities, false); } /* (non-Javadoc) diff --git a/src/main/java/org/springframework/data/cassandra/util/CqlUtils.java b/src/main/java/org/springframework/data/cassandra/util/CqlUtils.java index 28c89eae3..aee05dc7e 100644 --- a/src/main/java/org/springframework/data/cassandra/util/CqlUtils.java +++ b/src/main/java/org/springframework/data/cassandra/util/CqlUtils.java @@ -16,6 +16,7 @@ import com.datastax.driver.core.ColumnMetadata; import com.datastax.driver.core.DataType; import com.datastax.driver.core.Query; import com.datastax.driver.core.TableMetadata; +import com.datastax.driver.core.querybuilder.Batch; import com.datastax.driver.core.querybuilder.Delete; import com.datastax.driver.core.querybuilder.Delete.Where; import com.datastax.driver.core.querybuilder.Insert; @@ -245,6 +246,39 @@ public abstract class CqlUtils { } + /** + * Generates a Batch Object for multiple inserts + * + * @param keyspaceName + * @param tableName + * @param entity + * @param objectsToSave + * @param mappingContext + * @param beanClassLoader + * + * @return The Query object to run with session.execute(); + * @throws EntityWriterException + */ + public static Batch toInsertBatchQuery(final String keyspaceName, final String tableName, + final List objectsToSave, CassandraPersistentEntity entity) throws EntityWriterException { + + /* + * Return variable is a Batch statement + */ + final Batch b = QueryBuilder.batch(); + + List queries = new ArrayList(); + + for (final T objectToSave : objectsToSave) { + + queries.add(toInsertQuery(keyspaceName, tableName, objectToSave, entity)); + + } + + return b; + + } + /** * @param keyspace * @param tableName @@ -343,6 +377,10 @@ public abstract class CqlUtils { return str.toString(); } + /** + * @param dataType + * @return + */ public static String toCQL(DataType dataType) { if (dataType.getTypeArguments().isEmpty()) { return dataType.getName().name(); @@ -376,4 +414,32 @@ public abstract class CqlUtils { return str.toString(); } + /** + * @param keyspace + * @param tableName + * @param entities + * @param cPEntity + * @return + * @throws EntityWriterException + */ + public static Batch toDeleteBatchQuery(String keyspaceName, String tableName, List entities, + CassandraPersistentEntity entity) throws EntityWriterException { + + /* + * Return variable is a Batch statement + */ + final Batch b = QueryBuilder.batch(); + + List queries = new ArrayList(); + + for (final T objectToSave : entities) { + + queries.add(toDeleteQuery(keyspaceName, tableName, objectToSave, entity)); + + } + + return b; + + } + }