diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/MappingCassandraConverter.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/MappingCassandraConverter.java index fa69cf6b7..a2a74cff5 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/MappingCassandraConverter.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/core/convert/MappingCassandraConverter.java @@ -23,6 +23,7 @@ import java.util.ArrayList; import java.util.Collection; import java.util.Collections; import java.util.List; +import java.util.Map; import java.util.Optional; import org.slf4j.Logger; @@ -299,6 +300,7 @@ public class MappingCassandraConverter extends AbstractCassandraConverter } @Override + @SuppressWarnings("unchecked") public void write(Object source, Object sink, CassandraPersistentEntity entity) { if (source == null) { @@ -309,7 +311,9 @@ public class MappingCassandraConverter extends AbstractCassandraConverter throw new MappingException("No mapping metadata found for " + source.getClass()); } - if (sink instanceof Insert) { + if (sink instanceof Map) { + writeMapFromWrapper(getConvertingAccessor(source, entity), (Map) sink, entity); + } else if (sink instanceof Insert) { writeInsertFromObject(source, (Insert) sink, entity); } else if (sink instanceof Update) { writeUpdateFromObject(source, (Update) sink, entity); @@ -328,6 +332,36 @@ public class MappingCassandraConverter extends AbstractCassandraConverter writeInsertFromWrapper(getConvertingAccessor(object, entity), insert, entity); } + private void writeMapFromWrapper(final ConvertingPropertyAccessor accessor, final Map insert, + CassandraPersistentEntity entity) { + + entity.getPersistentProperties().forEach(property -> { + + Optional value = getWriteValue(property, accessor); + + if (log.isDebugEnabled()) { + log.debug("doWithProperties Property.type {}, Property.value {}", property.getType().getName(), value); + } + + if (property.isCompositePrimaryKey()) { + if (log.isDebugEnabled()) { + log.debug("Property is a compositeKey"); + } + + writeMapFromWrapper(getConvertingAccessor(value.orElse(null), property.getCompositePrimaryKeyEntity()), insert, + property.getCompositePrimaryKeyEntity()); + + return; + } + + if (log.isDebugEnabled()) { + log.debug("Adding map.entry [{}] - [{}]", property.getColumnName().toCql(), value); + } + + insert.put(property.getColumnName().toCql(), value.orElse(null)); + }); + } + protected void writeInsertFromWrapper(final ConvertingPropertyAccessor accessor, final Insert insert, CassandraPersistentEntity entity) { diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/CassandraEntityInformation.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/CassandraEntityInformation.java index 3bd2e8bc2..dc827d2f6 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/CassandraEntityInformation.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/query/CassandraEntityInformation.java @@ -25,13 +25,4 @@ import org.springframework.data.repository.core.EntityInformation; */ public interface CassandraEntityInformation extends EntityInformation, CassandraEntityMetadata { - /** - * Return {@literal true} if the persistent entity consists entirely of primary key properties (a single Id property, - * composite primary key). - * - * @return {@literal true} if the persistent entity consists entirely of primary key properties (a single Id property, - * composite primary key). - * @since 2.0 - */ - boolean isPrimaryKeyEntity(); } diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/MappingCassandraEntityInformation.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/MappingCassandraEntityInformation.java index 870d92faa..f9fd70df2 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/MappingCassandraEntityInformation.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/MappingCassandraEntityInformation.java @@ -16,7 +16,6 @@ package org.springframework.data.cassandra.repository.support; import java.util.Optional; -import java.util.concurrent.atomic.AtomicReference; import org.springframework.data.cassandra.core.convert.CassandraConverter; import org.springframework.data.cassandra.core.mapping.CassandraPersistentEntity; @@ -24,7 +23,6 @@ import org.springframework.data.cassandra.core.mapping.CassandraPersistentProper import org.springframework.data.cassandra.repository.MapId; import org.springframework.data.cassandra.repository.query.CassandraEntityInformation; import org.springframework.data.cql.core.CqlIdentifier; -import org.springframework.data.mapping.PropertyHandler; import org.springframework.data.repository.core.support.AbstractEntityInformation; import org.springframework.util.Assert; @@ -41,7 +39,6 @@ public class MappingCassandraEntityInformation extends AbstractEntityInfo private final CassandraPersistentEntity entityMetadata; private final CassandraConverter converter; - private final boolean isPrimaryKeyEntity; /** * Create a new {@link MappingCassandraEntityInformation} for the given {@link CassandraPersistentEntity}. @@ -54,7 +51,6 @@ public class MappingCassandraEntityInformation extends AbstractEntityInfo this.entityMetadata = entity; this.converter = converter; - this.isPrimaryKeyEntity = hasNonIdProperties(entity); } /* (non-Javadoc) @@ -89,28 +85,4 @@ public class MappingCassandraEntityInformation extends AbstractEntityInfo public CqlIdentifier getTableName() { return entityMetadata.getTableName(); } - - @Override - public boolean isPrimaryKeyEntity() { - return isPrimaryKeyEntity; - } - - private static boolean hasNonIdProperties(CassandraPersistentEntity entity) { - - final AtomicReference hasPrimaryKeyOnlyProperties = new AtomicReference(true); - - entity.doWithProperties(new PropertyHandler() { - @Override - public void doWithPersistentProperty(CassandraPersistentProperty property) { - - if (property.isCompositePrimaryKey() || property.isPrimaryKeyColumn() || property.isIdProperty()) { - return; - } - - hasPrimaryKeyOnlyProperties.set(false); - } - }); - - return hasPrimaryKeyOnlyProperties.get(); - } } 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 a2224f98a..553cf93ac 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 @@ -16,15 +16,21 @@ package org.springframework.data.cassandra.repository.support; import java.util.ArrayList; +import java.util.LinkedHashMap; import java.util.List; +import java.util.Map; +import java.util.Map.Entry; import java.util.Optional; import org.springframework.data.cassandra.core.CassandraOperations; import org.springframework.data.cassandra.core.CassandraTemplate; +import org.springframework.data.cassandra.core.convert.CassandraConverter; +import org.springframework.data.cassandra.core.mapping.CassandraPersistentEntity; import org.springframework.data.cassandra.repository.TypedIdCassandraRepository; import org.springframework.data.cassandra.repository.query.CassandraEntityInformation; import org.springframework.util.Assert; +import com.datastax.driver.core.querybuilder.Insert; import com.datastax.driver.core.querybuilder.QueryBuilder; import com.datastax.driver.core.querybuilder.Select; @@ -41,8 +47,6 @@ public class SimpleCassandraRepository implements TypedIdCassandraReposit private CassandraOperations operations; - private final boolean isPrimaryKeyEntity; - /** * Create a new {@link SimpleCassandraRepository} for the given {@link CassandraEntityInformation} and * {@link CassandraTemplate}. @@ -57,7 +61,6 @@ public class SimpleCassandraRepository implements TypedIdCassandraReposit this.entityInformation = metadata; this.operations = operations; - this.isPrimaryKeyEntity = metadata.isPrimaryKeyEntity(); } /* (non-Javadoc) @@ -68,39 +71,50 @@ public class SimpleCassandraRepository implements TypedIdCassandraReposit Assert.notNull(entity, "Entity must not be null"); - if (entityInformation.isNew(entity) || isPrimaryKeyEntity) { - return operations.insert(entity); - } + Insert insert = createFullInsert(entity); - return operations.update(entity); + operations.getCqlOperations().execute(insert); + + return entity; } /* (non-Javadoc) - * @see org.springframework.data.repository.CrudRepository#save(java.lang.Iterable) + * @see org.springframework.data.repository.CrudRepository#saveAll(java.lang.Iterable) */ @Override public List saveAll(Iterable entities) { Assert.notNull(entities, "The given Iterable of entities must not be null"); - List result = new ArrayList<>(); + + List result = new ArrayList(); + for (S entity : entities) { - S saved; - - if (entityInformation.isNew(entity) || isPrimaryKeyEntity) { - saved = operations.insert(entity); - } else { - saved = operations.update(entity); - } - - if (saved != null) { - result.add(saved); - } + result.add(entity); + operations.getCqlOperations().execute(createFullInsert(entity)); } return result; } + private Insert createFullInsert(S entity) { + + CassandraConverter converter = operations.getConverter(); + CassandraPersistentEntity persistentEntity = converter.getMappingContext() + .getRequiredPersistentEntity(entity.getClass()); + Map toInsert = new LinkedHashMap<>(); + + converter.write(entity, toInsert, persistentEntity); + + Insert insert = QueryBuilder.insertInto(persistentEntity.getTableName().toCql()); + + for (Entry entry : toInsert.entrySet()) { + insert.value(entry.getKey(), entry.getValue()); + } + + return insert; + } + /* (non-Javadoc) * @see org.springframework.data.cassandra.repository.TypedIdCassandraRepository#insert(java.lang.Object) */ diff --git a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/SimpleReactiveCassandraRepository.java b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/SimpleReactiveCassandraRepository.java index b83039d38..dd093951a 100644 --- a/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/SimpleReactiveCassandraRepository.java +++ b/spring-data-cassandra/src/main/java/org/springframework/data/cassandra/repository/support/SimpleReactiveCassandraRepository.java @@ -18,12 +18,19 @@ package org.springframework.data.cassandra.repository.support; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; +import java.util.LinkedHashMap; +import java.util.Map; +import java.util.Map.Entry; + import org.reactivestreams.Publisher; import org.springframework.data.cassandra.core.ReactiveCassandraOperations; +import org.springframework.data.cassandra.core.convert.CassandraConverter; +import org.springframework.data.cassandra.core.mapping.CassandraPersistentEntity; import org.springframework.data.cassandra.repository.ReactiveCassandraRepository; import org.springframework.data.cassandra.repository.query.CassandraEntityInformation; import org.springframework.util.Assert; +import com.datastax.driver.core.querybuilder.Insert; import com.datastax.driver.core.querybuilder.QueryBuilder; import com.datastax.driver.core.querybuilder.Select; @@ -39,8 +46,6 @@ public class SimpleReactiveCassandraRepository implements ReactiveCassand protected ReactiveCassandraOperations operations; - private final boolean isPrimaryKeyEntity; - /** * Create a new {@link SimpleReactiveCassandraRepository} for the given {@link CassandraEntityInformation} and * {@link ReactiveCassandraOperations}. @@ -56,7 +61,6 @@ public class SimpleReactiveCassandraRepository implements ReactiveCassand this.entityInformation = metadata; this.operations = operations; - this.isPrimaryKeyEntity = metadata.isPrimaryKeyEntity(); } /* (non-Javadoc) @@ -67,11 +71,7 @@ public class SimpleReactiveCassandraRepository implements ReactiveCassand Assert.notNull(entity, "Entity must not be null"); - if (entityInformation.isNew(entity) || isPrimaryKeyEntity) { - return operations.insert(entity); - } - - return operations.update(entity); + return operations.getReactiveCqlOperations().execute(createFullInsert(entity)).map(it -> entity); } @@ -95,15 +95,28 @@ public class SimpleReactiveCassandraRepository implements ReactiveCassand Assert.notNull(entityStream, "The given Publisher of entities must not be null"); return Flux.from(entityStream).flatMap(entity -> { - - if (entityInformation.isNew(entity) || isPrimaryKeyEntity) { - return operations.insert(entity); - } - - return operations.update(entity); + return operations.getReactiveCqlOperations().execute(createFullInsert(entity)).map(it -> entity); }); } + private Insert createFullInsert(S entity) { + + CassandraConverter converter = operations.getConverter(); + CassandraPersistentEntity persistentEntity = converter.getMappingContext() + .getRequiredPersistentEntity(entity.getClass()); + Map toInsert = new LinkedHashMap<>(); + + converter.write(entity, toInsert, persistentEntity); + + Insert insert = QueryBuilder.insertInto(persistentEntity.getTableName().toCql()); + + for (Entry entry : toInsert.entrySet()) { + insert.value(entry.getKey(), entry.getValue()); + } + + return insert; + } + /* (non-Javadoc) * @see org.springframework.data.cassandra.repository.ReactiveCassandraRepository#insert(java.lang.Object) */ diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/support/MappingCassandraEntityInformationUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/support/MappingCassandraEntityInformationUnitTests.java deleted file mode 100644 index b3b1cdbd1..000000000 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/support/MappingCassandraEntityInformationUnitTests.java +++ /dev/null @@ -1,106 +0,0 @@ -/* - * Copyright 2017 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ -package org.springframework.data.cassandra.repository.support; - -import static org.assertj.core.api.Assertions.*; - -import lombok.Data; - -import org.junit.Before; -import org.junit.Test; -import org.junit.runner.RunWith; -import org.mockito.Mock; -import org.mockito.junit.MockitoJUnitRunner; -import org.springframework.data.annotation.Id; -import org.springframework.data.cassandra.core.convert.CassandraConverter; -import org.springframework.data.cassandra.core.convert.MappingCassandraConverter; -import org.springframework.data.cassandra.core.mapping.BasicCassandraMappingContext; -import org.springframework.data.cassandra.core.mapping.UserTypeResolver; -import org.springframework.data.cassandra.domain.AllPossibleTypes; -import org.springframework.data.cassandra.domain.CompositeKey; -import org.springframework.data.cassandra.domain.TypeWithKeyClass; -import org.springframework.data.cassandra.domain.TypeWithMapId; - -/** - * Unit tests for {@link MappingCassandraEntityInformation}. - * - * @author Mark Paluch - */ -@RunWith(MockitoJUnitRunner.class) -@SuppressWarnings("unchecked") -public class MappingCassandraEntityInformationUnitTests { - - BasicCassandraMappingContext context = new BasicCassandraMappingContext(); - - CassandraConverter converter = new MappingCassandraConverter(context); - - @Mock UserTypeResolver userTypeResolver; - - @Before - public void before() { - context.setUserTypeResolver(userTypeResolver); - } - - @Test // DATACASS-420 - public void shouldConsiderSimpleIdEntityAsPrimaryKeyOnly() { - - MappingCassandraEntityInformation information = new MappingCassandraEntityInformation( - context.getRequiredPersistentEntity(PrimaryKeyOnly.class), converter); - - assertThat(information.isPrimaryKeyEntity()).isTrue(); - } - - @Test // DATACASS-420 - public void shouldConsiderCompositeIdEntityAsPrimaryKeyOnly() { - - MappingCassandraEntityInformation information = new MappingCassandraEntityInformation( - context.getRequiredPersistentEntity(CompositeKey.class), converter); - - assertThat(information.isPrimaryKeyEntity()).isTrue(); - } - - @Test // DATACASS-420 - public void shouldConsiderCompositeKeyClassEntityAsPrimaryKeyOnly() { - - MappingCassandraEntityInformation information = new MappingCassandraEntityInformation( - context.getRequiredPersistentEntity(TypeWithKeyClass.class), converter); - - assertThat(information.isPrimaryKeyEntity()).isTrue(); - } - - @Test // DATACASS-420 - public void shouldConsiderMapIdClassEntityAsPrimaryKeyOnly() { - - MappingCassandraEntityInformation information = new MappingCassandraEntityInformation( - context.getRequiredPersistentEntity(TypeWithMapId.class), converter); - - assertThat(information.isPrimaryKeyEntity()).isTrue(); - } - - @Test // DATACASS-420 - public void shouldComplexEntityNotAsPrimaryKeyOnly() { - - MappingCassandraEntityInformation information = new MappingCassandraEntityInformation( - context.getRequiredPersistentEntity(AllPossibleTypes.class), converter); - - assertThat(information.isPrimaryKeyEntity()).isFalse(); - } - - @Data - static class PrimaryKeyOnly { - @Id String id; - } -} diff --git a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/support/SimpleCassandraRepositoryUnitTests.java b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/support/SimpleCassandraRepositoryUnitTests.java index 2d95683ba..7fa4967eb 100644 --- a/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/support/SimpleCassandraRepositoryUnitTests.java +++ b/spring-data-cassandra/src/test/java/org/springframework/data/cassandra/repository/support/SimpleCassandraRepositoryUnitTests.java @@ -25,6 +25,8 @@ import java.io.Serializable; import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Captor; import org.mockito.Mock; import org.mockito.junit.MockitoJUnitRunner; import org.springframework.data.annotation.Id; @@ -34,6 +36,11 @@ import org.springframework.data.cassandra.core.mapping.BasicCassandraMappingCont import org.springframework.data.cassandra.core.mapping.CassandraPersistentEntity; import org.springframework.data.cassandra.core.mapping.UserTypeResolver; import org.springframework.data.cassandra.domain.Person; +import org.springframework.data.cql.core.CqlIdentifier; +import org.springframework.data.cql.core.CqlOperations; + +import com.datastax.driver.core.UserType; +import com.datastax.driver.core.querybuilder.Insert; /** * Unit tests for {@link SimpleCassandraRepository}. @@ -50,11 +57,20 @@ public class SimpleCassandraRepositoryUnitTests { SimpleCassandraRepository repository; @Mock CassandraOperations cassandraOperations; + @Mock CqlOperations cqlOperations; @Mock UserTypeResolver userTypeResolver; + @Mock UserType userType; + + @Captor ArgumentCaptor insertCaptor; @Before public void before() { + mappingContext.setUserTypeResolver(userTypeResolver); + + when(cassandraOperations.getConverter()).thenReturn(converter); + when(cassandraOperations.getCqlOperations()).thenReturn(cqlOperations); + when(userTypeResolver.resolveType(CqlIdentifier.cqlId("address"))).thenReturn(userType); } @Test // DATACASS-428 @@ -67,12 +83,10 @@ public class SimpleCassandraRepositoryUnitTests { SimplePerson person = new SimplePerson(); - when(cassandraOperations.insert(person)).thenReturn(person); + repository.save(person); - Object result = repository.save(person); - - assertThat(result).isEqualTo(person); - verify(cassandraOperations).insert(person); + verify(cqlOperations).execute(insertCaptor.capture()); + assertThat(insertCaptor.getValue().toString()).isEqualTo("INSERT INTO simpleperson (id) VALUES (null);"); } @Test // DATACASS-428 @@ -85,12 +99,9 @@ public class SimpleCassandraRepositoryUnitTests { Person person = new Person(); - when(cassandraOperations.update(person)).thenReturn(person); + repository.save(person); - Object result = repository.save(person); - - assertThat(result).isEqualTo(person); - verify(cassandraOperations).update(person); + verify(cqlOperations).execute(any(Insert.class)); } @Test // DATACASS-428 @@ -105,12 +116,12 @@ public class SimpleCassandraRepositoryUnitTests { person.setFirstname("foo"); person.setLastname("bar"); - when(cassandraOperations.update(person)).thenReturn(person); + repository.save(person); - Object result = repository.save(person); - - assertThat(result).isEqualTo(person); - verify(cassandraOperations).update(person); + verify(cqlOperations).execute(insertCaptor.capture()); + assertThat(insertCaptor.getValue().toString()) + .contains("INSERT INTO person (lastname,firstname,alternativeaddresses,"); + assertThat(insertCaptor.getValue().toString()).contains("VALUES ('bar','foo',null"); } @Test // DATACASS-428 @@ -123,11 +134,8 @@ public class SimpleCassandraRepositoryUnitTests { Person person = new Person(); - when(cassandraOperations.insert(person)).thenReturn(person); + repository.insert(person); - Object result = repository.insert(person); - - assertThat(result).isEqualTo(person); verify(cassandraOperations).insert(person); }