#12 - Retain input item order in save(…) methods.

We now use concatMap(…) instead of flatMap(…) when saving items to retain output item order.
This commit is contained in:
Mark Paluch
2018-11-12 15:13:19 +01:00
parent 93ddaff3c2
commit 226c80fb02
2 changed files with 26 additions and 22 deletions

View File

@@ -15,16 +15,13 @@
*/
package org.springframework.data.r2dbc.repository.support;
import lombok.NonNull;
import lombok.RequiredArgsConstructor;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;
import java.util.stream.IntStream;
import lombok.NonNull;
import lombok.RequiredArgsConstructor;
import org.reactivestreams.Publisher;
import org.springframework.data.r2dbc.function.DatabaseClient;
import org.springframework.data.r2dbc.function.DatabaseClient.BindSpec;
@@ -35,6 +32,8 @@ import org.springframework.data.r2dbc.function.convert.SettableValue;
import org.springframework.data.relational.repository.query.RelationalEntityInformation;
import org.springframework.data.repository.reactive.ReactiveCrudRepository;
import org.springframework.util.Assert;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
/**
* Simple {@link ReactiveCrudRepository} implementation using R2DBC through {@link DatabaseClient}.
@@ -118,7 +117,7 @@ public class SimpleR2dbcRepository<T, ID> implements ReactiveCrudRepository<T, I
Assert.notNull(objectsToSave, "Objects to save must not be null!");
return Flux.fromIterable(objectsToSave).flatMap(this::save);
return Flux.fromIterable(objectsToSave).concatMap(this::save);
}
/* (non-Javadoc)
@@ -129,7 +128,7 @@ public class SimpleR2dbcRepository<T, ID> implements ReactiveCrudRepository<T, I
Assert.notNull(objectsToSave, "Object publisher must not be null!");
return Flux.from(objectsToSave).flatMap(this::save);
return Flux.from(objectsToSave).concatMap(this::save);
}
/* (non-Javadoc)
@@ -210,7 +209,7 @@ public class SimpleR2dbcRepository<T, ID> implements ReactiveCrudRepository<T, I
Assert.notNull(idPublisher, "The Id Publisher must not be null!");
return Flux.from(idPublisher).buffer().filter(ids -> !ids.isEmpty()).flatMap(ids -> {
return Flux.from(idPublisher).buffer().filter(ids -> !ids.isEmpty()).concatMap(ids -> {
String bindings = getInBinding(ids);
@@ -259,7 +258,7 @@ public class SimpleR2dbcRepository<T, ID> implements ReactiveCrudRepository<T, I
Assert.notNull(idPublisher, "The Id Publisher must not be null!");
return Flux.from(idPublisher).buffer().filter(ids -> !ids.isEmpty()).flatMap(ids -> {
return Flux.from(idPublisher).buffer().filter(ids -> !ids.isEmpty()).concatMap(ids -> {
String bindings = getInBinding(ids);

View File

@@ -17,26 +17,21 @@ package org.springframework.data.r2dbc.repository.support;
import static org.assertj.core.api.Assertions.*;
import io.r2dbc.spi.ConnectionFactory;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Hooks;
import reactor.core.publisher.Mono;
import reactor.test.StepVerifier;
import java.util.Arrays;
import java.util.Collections;
import java.util.Map;
import io.r2dbc.spi.ConnectionFactory;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.junit.Before;
import org.junit.Test;
import org.springframework.data.annotation.Id;
import org.springframework.data.r2dbc.testing.R2dbcIntegrationTestSupport;
import org.springframework.data.r2dbc.function.DatabaseClient;
import org.springframework.data.r2dbc.function.DefaultReactiveDataAccessStrategy;
import org.springframework.data.r2dbc.function.convert.MappingR2dbcConverter;
import org.springframework.data.r2dbc.testing.R2dbcIntegrationTestSupport;
import org.springframework.data.relational.core.conversion.BasicRelationalConverter;
import org.springframework.data.relational.core.mapping.RelationalMappingContext;
import org.springframework.data.relational.core.mapping.RelationalPersistentEntity;
@@ -44,6 +39,10 @@ import org.springframework.data.relational.core.mapping.Table;
import org.springframework.data.relational.repository.query.RelationalEntityInformation;
import org.springframework.data.relational.repository.support.MappingRelationalEntityInformation;
import org.springframework.jdbc.core.JdbcTemplate;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Hooks;
import reactor.core.publisher.Mono;
import reactor.test.StepVerifier;
/**
* Integration tests for {@link SimpleR2dbcRepository}.
@@ -122,14 +121,20 @@ public class SimpleR2dbcRepositoryIntegrationTests extends R2dbcIntegrationTestS
LegoSet legoSet1 = new LegoSet(null, "SCHAUFELRADBAGGER", 12);
LegoSet legoSet2 = new LegoSet(null, "FORSCHUNGSSCHIFF", 13);
LegoSet legoSet3 = new LegoSet(null, "RALLYEAUTO", 14);
LegoSet legoSet4 = new LegoSet(null, "VOLTRON", 15);
repository.saveAll(Arrays.asList(legoSet1, legoSet2)) //
repository.saveAll(Arrays.asList(legoSet1, legoSet2, legoSet3, legoSet4)) //
.map(LegoSet::getManual) //
.as(StepVerifier::create) //
.expectNextCount(2) //
.expectNext(12) //
.expectNext(13) //
.expectNext(14) //
.expectNext(15) //
.verifyComplete();
Map<String, Object> map = jdbc.queryForMap("SELECT COUNT(*) FROM repo_legoset");
assertThat(map).containsEntry("count", 2L);
assertThat(map).containsEntry("count", 4L);
}
@Test