From 54d5a529e4bb4aaf580a77919f367a80ca0a5098 Mon Sep 17 00:00:00 2001 From: Michael Simons Date: Tue, 24 Oct 2023 10:26:00 +0200 Subject: [PATCH] refactor: Avoid concurrency issues during some writes and reads. --- .../data/neo4j/core/ReactiveNeo4jTemplate.java | 12 ++++++------ .../support/SimpleReactiveNeo4jRepository.java | 6 +++--- 2 files changed, 9 insertions(+), 9 deletions(-) diff --git a/src/main/java/org/springframework/data/neo4j/core/ReactiveNeo4jTemplate.java b/src/main/java/org/springframework/data/neo4j/core/ReactiveNeo4jTemplate.java index 99f2cb661..259caeb9e 100644 --- a/src/main/java/org/springframework/data/neo4j/core/ReactiveNeo4jTemplate.java +++ b/src/main/java/org/springframework/data/neo4j/core/ReactiveNeo4jTemplate.java @@ -401,7 +401,7 @@ public final class ReactiveNeo4jTemplate implements NestedRelationshipProcessingStateMachine stateMachine = new NestedRelationshipProcessingStateMachine(neo4jMappingContext); EntityFromDtoInstantiatingConverter converter = new EntityFromDtoInstantiatingConverter<>(domainType, neo4jMappingContext); return Flux.fromIterable(instances) - .flatMap(instance -> { + .concatMap(instance -> { T domainObject = converter.convert(instance); @SuppressWarnings("unchecked") @@ -534,7 +534,7 @@ public final class ReactiveNeo4jTemplate implements Neo4jPersistentEntity entityMetaData = neo4jMappingContext.getRequiredPersistentEntity(commonElementType); Neo4jPersistentProperty idProperty = entityMetaData.getIdProperty(); - return savedInstances.flatMap(savedInstance -> { + return savedInstances.concatMap(savedInstance -> { PersistentPropertyAccessor propertyAccessor = entityMetaData.getPropertyAccessor(savedInstance); return findById(propertyAccessor.getProperty(idProperty), commonElementType); }).map(instance -> localProjectionFactory.createProjection(resultType, instance)); @@ -595,7 +595,7 @@ public final class ReactiveNeo4jTemplate implements .all() .collectMap(m -> (Value) m.getT1(), m -> (String) m.getT2()); }).flatMapMany(idToInternalIdMapping -> Flux.fromIterable(entitiesToBeSaved) - .flatMap(t -> { + .concatMap(t -> { PersistentPropertyAccessor propertyAccessor = entityMetaData.getPropertyAccessor(t.getT3()); Neo4jPersistentProperty idProperty = entityMetaData.getRequiredIdProperty(); Object id = convertIdValues(idProperty, propertyAccessor.getProperty(idProperty)); @@ -722,7 +722,7 @@ public final class ReactiveNeo4jTemplate implements Set processedRelationshipIds = ctx.get("processedRelationships"); Set processedNodeIds = ctx.get("processedNodes"); return Flux.fromIterable(entityMetaData.getRelationshipsInHierarchy(queryFragments::includeField)) - .flatMap(relationshipDescription -> { + .concatMap(relationshipDescription -> { Statement statement = cypherGenerator.prepareMatchOf(entityMetaData, relationshipDescription, queryFragments.getMatchOn(), queryFragments.getCondition()) @@ -777,7 +777,7 @@ public final class ReactiveNeo4jTemplate implements return queryFragments.includeField(prepend); } )) - .flatMap(relDe -> { + .concatMap(relDe -> { Node node = anyNode(Constants.NAME_OF_TYPED_ROOT_NODE.apply(target)); Statement statement = cypherGenerator @@ -1188,7 +1188,7 @@ public final class ReactiveNeo4jTemplate implements return fetchSpec.all().switchOnFirst((signal, f) -> { if (signal.hasValue() && preparedQuery.resultsHaveBeenAggregated()) { - return f.flatMap(nested -> Flux.fromIterable((Collection) nested).distinct()).distinct(); + return f.concatMap(nested -> Flux.fromIterable((Collection) nested).distinct()).distinct(); } return f; }); diff --git a/src/main/java/org/springframework/data/neo4j/repository/support/SimpleReactiveNeo4jRepository.java b/src/main/java/org/springframework/data/neo4j/repository/support/SimpleReactiveNeo4jRepository.java index 989f903ec..22a1085b7 100644 --- a/src/main/java/org/springframework/data/neo4j/repository/support/SimpleReactiveNeo4jRepository.java +++ b/src/main/java/org/springframework/data/neo4j/repository/support/SimpleReactiveNeo4jRepository.java @@ -85,7 +85,7 @@ public class SimpleReactiveNeo4jRepository implements ReactiveSortingRepo @Override public Flux findAllById(Publisher idStream) { - return Flux.from(idStream).buffer().flatMap(this::findAllById); + return Flux.from(idStream).buffer().concatMap(this::findAllById); } @Override @@ -135,7 +135,7 @@ public class SimpleReactiveNeo4jRepository implements ReactiveSortingRepo @Transactional public Flux saveAll(Publisher entityStream) { - return Flux.from(entityStream).flatMap(this::save); + return Flux.from(entityStream).concatMap(this::save); } /* @@ -218,7 +218,7 @@ public class SimpleReactiveNeo4jRepository implements ReactiveSortingRepo public Mono deleteAll(Publisher entitiesPublisher) { Assert.notNull(entitiesPublisher, "The given Publisher of entities must not be null"); - return Flux.from(entitiesPublisher).flatMap(this::delete).then(); + return Flux.from(entitiesPublisher).concatMap(this::delete).then(); } /*