From 42a40bb2b0db417b9f38ee80556ebb51098e0dc3 Mon Sep 17 00:00:00 2001 From: dwebb Date: Tue, 19 Nov 2013 16:24:01 -0500 Subject: [PATCH] DATACASS-32 - Pulled in existing work from Matthew Adams on the CassandraTemplate --- .../cassandra/core/CassandraOperations.java | 20 +++ .../cassandra/core/CassandraTemplate.java | 154 ++++++++++++++++++ .../cassandra/core/ResultSetExtractor.java | 11 ++ .../cassandra/core/RowCallbackHandler.java | 10 ++ .../cassandra/core/RowMapper.java | 10 ++ 5 files changed, 205 insertions(+) create mode 100644 src/main/java/org/springframework/cassandra/core/ResultSetExtractor.java create mode 100644 src/main/java/org/springframework/cassandra/core/RowCallbackHandler.java create mode 100644 src/main/java/org/springframework/cassandra/core/RowMapper.java diff --git a/src/main/java/org/springframework/cassandra/core/CassandraOperations.java b/src/main/java/org/springframework/cassandra/core/CassandraOperations.java index f3f7a7e1e..730f0b078 100644 --- a/src/main/java/org/springframework/cassandra/core/CassandraOperations.java +++ b/src/main/java/org/springframework/cassandra/core/CassandraOperations.java @@ -15,6 +15,10 @@ */ package org.springframework.cassandra.core; +import java.util.List; +import java.util.Map; + +import org.springframework.dao.DataAccessException; import com.datastax.driver.core.ResultSet; import com.datastax.driver.core.ResultSetFuture; @@ -36,6 +40,22 @@ public interface CassandraOperations { RuntimeException potentiallyConvertRuntimeException(RuntimeException ex); + T query(String cql, ResultSetExtractor rse) throws DataAccessException; + + void query(String cql, RowCallbackHandler rch) throws DataAccessException; + + List query(String cql, RowMapper rowMapper) throws DataAccessException; + + T queryForObject(String cql, RowMapper rowMapper) throws DataAccessException; + + T queryForObject(String cql, Class requiredType) throws DataAccessException; + + Map queryForMap(String cql) throws DataAccessException; + + List queryForList(String cql, Class elementType) throws DataAccessException; + + List> queryForList(String cql) throws DataAccessException; + /** * Execute query and return Cassandra ResultSet * diff --git a/src/main/java/org/springframework/cassandra/core/CassandraTemplate.java b/src/main/java/org/springframework/cassandra/core/CassandraTemplate.java index 8c763f6e7..36122f9f8 100644 --- a/src/main/java/org/springframework/cassandra/core/CassandraTemplate.java +++ b/src/main/java/org/springframework/cassandra/core/CassandraTemplate.java @@ -15,14 +15,23 @@ */ package org.springframework.cassandra.core; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + import org.springframework.cassandra.support.CassandraAccessor; import org.springframework.dao.DataAccessException; import org.springframework.data.cassandra.core.CassandraDataTemplate; import org.springframework.util.Assert; +import com.datastax.driver.core.ColumnDefinitions; +import com.datastax.driver.core.ColumnDefinitions.Definition; 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.DriverException; /** * The CassandraTemplate is a Spring convenience wrapper for low level and explicit operations on the Cassandra @@ -142,4 +151,149 @@ public class CassandraTemplate extends CassandraAccessor implements CassandraOpe } } + /* (non-Javadoc) + * @see org.springframework.cassandra.core.CassandraOperations#query(java.lang.String, org.springframework.cassandra.core.ResultSetExtractor) + */ + public T query(String cql, ResultSetExtractor rse) throws DataAccessException { + try { + ResultSet rs = getSession().execute(cql); + return rse.extractData(rs); + } catch (DriverException dx) { + throw getExceptionTranslator().translateExceptionIfPossible(dx); + } + } + + /* (non-Javadoc) + * @see org.springframework.cassandra.core.CassandraOperations#query(java.lang.String, org.springframework.cassandra.core.RowCallbackHandler) + */ + public void query(String cql, RowCallbackHandler rch) throws DataAccessException { + try { + ResultSet rs = getSession().execute(cql); + for (Row row : rs.all()) { + rch.processRow(row); + } + } catch (DriverException dx) { + throw getExceptionTranslator().translateExceptionIfPossible(dx); + } + + } + + /* (non-Javadoc) + * @see org.springframework.cassandra.core.CassandraOperations#query(java.lang.String, org.springframework.cassandra.core.RowMapper) + */ + public List query(String cql, RowMapper rowMapper) throws DataAccessException { + try { + ResultSet rs = getSession().execute(cql); + int i = 0; + List mappedRows = new ArrayList(); + for (Row row : rs.all()) { + mappedRows.add(rowMapper.mapRow(row, i++)); + } + return mappedRows; + } catch (DriverException dx) { + throw getExceptionTranslator().translateExceptionIfPossible(dx); + } + } + + /* (non-Javadoc) + * @see org.springframework.cassandra.core.CassandraOperations#queryForObject(java.lang.String, org.springframework.cassandra.core.RowMapper) + */ + public T queryForObject(String cql, RowMapper rowMapper) throws DataAccessException { + try { + ResultSet rs = getSession().execute(cql); + List rows = rs.all(); + Assert.notNull(rows, "null row list returned from query"); + Assert.isTrue(rows.size() == 1, "row list has " + rows.size() + " rows instead of one"); + return rowMapper.mapRow(rows.get(0), 0); + } catch (DriverException dx) { + throw getExceptionTranslator().translateExceptionIfPossible(dx); + } + } + + /* (non-Javadoc) + * @see org.springframework.cassandra.core.CassandraOperations#queryForObject(java.lang.String, java.lang.Class) + */ + @SuppressWarnings("unchecked") + public T queryForObject(String cql, Class requiredType) throws DataAccessException { + ResultSet rs = getSession().execute(cql); + if (rs == null) { + return null; + } + Row row = rs.one(); + if (row == null) { + return null; + } + return (T) firstColumnToObject(row); + } + + /** + * @param row + * @return + */ + protected Object firstColumnToObject(Row row) { + ColumnDefinitions cols = row.getColumnDefinitions(); + if (cols.size() == 0) { + return null; + } + return cols.getType(0).deserialize(row.getBytesUnsafe(0)); + } + + /* (non-Javadoc) + * @see org.springframework.cassandra.core.CassandraOperations#queryForMap(java.lang.String) + */ + public Map queryForMap(String cql) throws DataAccessException { + ResultSet rs = getSession().execute(cql); + if (rs == null) { + return null; + } + return toMap(rs.one()); + } + + /** + * @param row + * @return + */ + protected Map toMap(Row row) { + if (row == null) { + return null; + } + + ColumnDefinitions cols = row.getColumnDefinitions(); + Map map = new HashMap(cols.size()); + + for (Definition def : cols.asList()) { + String name = def.getName(); + map.put(name, def.getType().deserialize(row.getBytesUnsafe(name))); + } + + return map; + } + + /* (non-Javadoc) + * @see org.springframework.cassandra.core.CassandraOperations#queryForList(java.lang.String, java.lang.Class) + */ + @SuppressWarnings("unchecked") + public List queryForList(String cql, Class elementType) throws DataAccessException { + ResultSet rs = getSession().execute(cql); + List rows = rs.all(); + List list = new ArrayList(rows.size()); + for (Row row : rows) { + list.add((T) firstColumnToObject(row)); + } + return list; + } + + /* (non-Javadoc) + * @see org.springframework.cassandra.core.CassandraOperations#queryForList(java.lang.String) + */ + public List> queryForList(String cql) throws DataAccessException { + ResultSet rs = getSession().execute(cql); + List rows = rs.all(); + List> list = new ArrayList>(rows.size()); + for (Row row : rows) { + list.add(toMap(row)); + } + return list; + } + } diff --git a/src/main/java/org/springframework/cassandra/core/ResultSetExtractor.java b/src/main/java/org/springframework/cassandra/core/ResultSetExtractor.java new file mode 100644 index 000000000..94b03aae9 --- /dev/null +++ b/src/main/java/org/springframework/cassandra/core/ResultSetExtractor.java @@ -0,0 +1,11 @@ +package org.springframework.cassandra.core; + +import org.springframework.dao.DataAccessException; + +import com.datastax.driver.core.ResultSet; +import com.datastax.driver.core.exceptions.DriverException; + +public interface ResultSetExtractor { + + T extractData(ResultSet rs) throws DriverException, DataAccessException; +} diff --git a/src/main/java/org/springframework/cassandra/core/RowCallbackHandler.java b/src/main/java/org/springframework/cassandra/core/RowCallbackHandler.java new file mode 100644 index 000000000..8eaf6da8e --- /dev/null +++ b/src/main/java/org/springframework/cassandra/core/RowCallbackHandler.java @@ -0,0 +1,10 @@ +package org.springframework.cassandra.core; + +import com.datastax.driver.core.Row; +import com.datastax.driver.core.exceptions.DriverException; + +public interface RowCallbackHandler { + + void processRow(Row row) throws DriverException; + +} diff --git a/src/main/java/org/springframework/cassandra/core/RowMapper.java b/src/main/java/org/springframework/cassandra/core/RowMapper.java new file mode 100644 index 000000000..4de62c128 --- /dev/null +++ b/src/main/java/org/springframework/cassandra/core/RowMapper.java @@ -0,0 +1,10 @@ +package org.springframework.cassandra.core; + +import com.datastax.driver.core.Row; +import com.datastax.driver.core.exceptions.DriverException; + +public interface RowMapper { + + T mapRow(Row row, int rowNum) throws DriverException; + +}