diff --git a/spring-cql/src/main/java/org/springframework/cassandra/core/CqlTemplate.java b/spring-cql/src/main/java/org/springframework/cassandra/core/CqlTemplate.java
index e9bf58586..c3b275fdc 100644
--- a/spring-cql/src/main/java/org/springframework/cassandra/core/CqlTemplate.java
+++ b/spring-cql/src/main/java/org/springframework/cassandra/core/CqlTemplate.java
@@ -52,6 +52,7 @@ import org.springframework.dao.DataAccessException;
import org.springframework.dao.IncorrectResultSizeDataAccessException;
import org.springframework.dao.InvalidDataAccessApiUsageException;
import org.springframework.dao.QueryTimeoutException;
+import org.springframework.dao.support.PersistenceExceptionTranslator;
import org.springframework.util.Assert;
import com.datastax.driver.core.BoundStatement;
@@ -854,19 +855,49 @@ public class CqlTemplate extends CassandraAccessor implements CqlOperations {
}
/**
- * Attempts to translate the {@link RuntimeException} into a Spring Data {@link Exception}.
+ * Attempts to translate the {@link Exception} into a Spring Data {@link Exception}.
+ * @param ex the Exception
+ * @return the translated {@link RuntimeException}
*/
@SuppressWarnings("all")
- protected RuntimeException translateExceptionIfPossible(RuntimeException e) {
-
- RuntimeException resolved = getExceptionTranslator().translateExceptionIfPossible(e);
- return (resolved != null ? resolved : e);
+ protected RuntimeException translateExceptionIfPossible(Exception ex) {
+ return translateExceptionIfPossible(ex, getExceptionTranslator());
}
+ /**
+ * Tries to convert the given {@link RuntimeException} into a {@link DataAccessException} but returns the original
+ * exception if the conversation failed. Thus allows safe re-throwing of the return value.
+ *
+ * @param ex the exception to translate
+ * @param exceptionTranslator the {@link PersistenceExceptionTranslator} to be used for translation
+ * @return
+ */
@SuppressWarnings("all")
- protected RuntimeException translateExceptionIfPossible(Exception e) {
- return (e instanceof RuntimeException ? translateExceptionIfPossible((RuntimeException) e)
- : new CassandraUncategorizedDataAccessException("Caught Uncategorized Exception", e));
+ protected static RuntimeException translateExceptionIfPossible(Exception ex, PersistenceExceptionTranslator exceptionTranslator) {
+
+ Assert.notNull(ex, "Exception must not be null");
+ Assert.notNull(exceptionTranslator, "PersistenceExceptionTranslator must not be null");
+
+ if (ex instanceof RuntimeException) {
+ return potentiallyConvertRuntimeException((RuntimeException) ex, exceptionTranslator);
+ }
+
+ return new CassandraUncategorizedDataAccessException("Caught Uncategorized Exception", ex);
+ }
+
+ /**
+ * Tries to convert the given {@link RuntimeException} into a {@link DataAccessException} but returns the original
+ * exception if the conversation failed. Thus allows safe re-throwing of the return value.
+ *
+ * @param ex the exception to translate
+ * @param exceptionTranslator the {@link PersistenceExceptionTranslator} to be used for translation
+ * @return
+ */
+ private static RuntimeException potentiallyConvertRuntimeException(RuntimeException ex,
+ PersistenceExceptionTranslator exceptionTranslator) {
+
+ RuntimeException resolved = exceptionTranslator.translateExceptionIfPossible(ex);
+ return resolved == null ? ex : resolved;
}
@Override
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 b306921ca..ea43ca596 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
@@ -15,6 +15,7 @@
*/
package org.springframework.data.cassandra.core;
+import java.util.Iterator;
import java.util.List;
import org.springframework.cassandra.core.Cancellable;
@@ -51,6 +52,20 @@ public interface CassandraOperations extends CqlOperations {
*/
CqlIdentifier getTableName(Class> entityClass);
+ /**
+ * Executes the given select {@code query} on the entity table of the specified {@code type} backed by a Cassandra
+ * {@link com.datastax.driver.core.ResultSet}.
+ *
+ * Returns a {@link java.util.Iterator} that wraps the a Cassandra {@link com.datastax.driver.core.ResultSet}.
+ *
+ * @param element return type
+ * @param query must not be empty and not {@literal null}.
+ * @param type must not be {@literal null}.
+ * @return
+ * @since 1.5
+ */
+ Iterator stream(String query, Class type);
+
/**
* Execute query and convert ResultSet to the list of entities.
*
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 a90c2cbf5..7440a4d88 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
@@ -17,6 +17,7 @@ package org.springframework.data.cassandra.core;
import java.util.ArrayList;
import java.util.Collection;
+import java.util.Collections;
import java.util.Iterator;
import java.util.List;
import java.util.Map;
@@ -34,6 +35,7 @@ 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.dao.support.PersistenceExceptionTranslator;
import org.springframework.data.cassandra.convert.CassandraConverter;
import org.springframework.data.cassandra.convert.MappingCassandraConverter;
import org.springframework.data.cassandra.mapping.CassandraMappingContext;
@@ -588,6 +590,29 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation
return result;
}
+ /* (non-Javadoc)
+ * @see org.springframework.data.cassandra.core.CassandraOperations#stream(java.lang.String, java.lang.Class)
+ */
+ public Iterator stream(final String query, Class type) {
+
+ Assert.hasText(query, "Query must not be empty");
+ Assert.notNull(type, "Type must not be null");
+
+ ResultSet resultSet = doExecute(new SessionCallback() {
+
+ @Override
+ public ResultSet doInSession(Session s) throws DataAccessException {
+ return s.execute(query);
+ }
+ });
+
+ if (resultSet == null) {
+ return Collections.emptyList().iterator();
+ }
+
+ return new ResultSetIteratorAdapter(resultSet.iterator(), getExceptionTranslator(), new CassandraConverterRowCallback(cassandraConverter, type));
+ }
+
protected List select(final Select query, CassandraConverterRowCallback readRowCallback) {
ResultSet resultSet = doExecute(new SessionCallback() {
@@ -1067,8 +1092,7 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation
throw new DuplicateKeyException("found two or more results in query " + query);
}
listener.onQueryComplete(result);
- }
- else{
+ } else {
listener.onQueryComplete(null);
}
} catch (Exception e) {
@@ -1085,4 +1109,38 @@ public class CassandraTemplate extends CqlTemplate implements CassandraOperation
throw new IllegalArgumentException(
String.format("Expected type String or Select; got type [%s] with value [%s]", query.getClass(), query));
}
+
+ private static class ResultSetIteratorAdapter implements Iterator{
+
+ private final Iterator iterator;
+ private final PersistenceExceptionTranslator exceptionTranslator;
+ private final CassandraConverterRowCallback rowCallback;
+
+ public ResultSetIteratorAdapter(Iterator iterator, PersistenceExceptionTranslator exceptionTranslator, CassandraConverterRowCallback rowCallback) {
+
+ this.iterator = iterator;
+ this.exceptionTranslator = exceptionTranslator;
+ this.rowCallback = rowCallback;
+ }
+
+ @Override
+ public boolean hasNext() {
+
+ try {
+ return iterator.hasNext();
+ } catch (Exception e) {
+ throw translateExceptionIfPossible(e, exceptionTranslator);
+ }
+ }
+
+ @Override
+ public T next() {
+
+ try {
+ return rowCallback.doWith(iterator.next());
+ } catch (Exception e) {
+ throw translateExceptionIfPossible(e, exceptionTranslator);
+ }
+ }
+ }
}
diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractCassandraQuery.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractCassandraQuery.java
index d9e447152..ccb89c6e6 100644
--- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractCassandraQuery.java
+++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/AbstractCassandraQuery.java
@@ -34,6 +34,7 @@ import org.springframework.data.cassandra.repository.query.CassandraQueryExecuti
import org.springframework.data.cassandra.repository.query.CassandraQueryExecution.ResultProcessingExecution;
import org.springframework.data.cassandra.repository.query.CassandraQueryExecution.ResultSetQuery;
import org.springframework.data.cassandra.repository.query.CassandraQueryExecution.SingleEntityExecution;
+import org.springframework.data.cassandra.repository.query.CassandraQueryExecution.StreamExecution;
import org.springframework.data.repository.query.ParameterAccessor;
import org.springframework.data.repository.query.RepositoryQuery;
import org.springframework.data.repository.query.ResultProcessor;
@@ -100,13 +101,15 @@ public abstract class AbstractCassandraQuery implements RepositoryQuery {
private CassandraQueryExecution getExecution(String query, CassandraParameterAccessor accessor,
Converter