From 6259cd2c3b00f0f58693b20652f556d5afa83b8f Mon Sep 17 00:00:00 2001 From: Christoph Strobl Date: Thu, 6 Feb 2020 08:16:59 +0100 Subject: [PATCH] DATAMONGO-2341 - Support shard key derivation in save operations via @Sharded annotation. Spring Data MongoDB uses the @Sharded annotation to identify entities stored in sharded collections. The shard key consists of a single or multiple properties present in every document within the target collection, and is used to distribute them across shards. Spring Data MongoDB will do best effort optimisations for sharded scenarios when using repositories by adding required shard key information, if not already present, to replaceOne filter queries when upserting entities. This may require an additional server round trip to determine the actual value of the current shard key. By setting @Sharded(immutableKey = true) no attempt will be made to check if an entities shard key changed. Please see the MongoDB Documentation for further details and the list below for which operations are eligible to auto include the shard key. * Reactive/CrudRepository.save(...) * Reactive/CrudRepository.saveAll(...) * Reactive/MongoTemplate.save(...) Original pull request: #833. --- .../data/mongodb/core/MappedDocument.java | 6 +- .../data/mongodb/core/MongoTemplate.java | 48 ++++-- .../data/mongodb/core/QueryOperations.java | 76 ++++++++- .../mongodb/core/ReactiveMongoTemplate.java | 44 ++++- .../mapping/BasicMongoPersistentEntity.java | 30 ++++ .../core/mapping/MongoPersistentEntity.java | 34 ++++ .../data/mongodb/core/mapping/ShardKey.java | 138 ++++++++++++++++ .../data/mongodb/core/mapping/Sharded.java | 92 +++++++++++ .../core/mapping/ShardingStrategy.java | 35 ++++ .../data/mongodb/util/BsonUtils.java | 4 +- .../mongodb/core/MongoTemplateUnitTests.java | 89 +++++++++++ .../core/ReactiveMongoTemplateUnitTests.java | 96 +++++++++++ .../ShardedEntityWithDefaultShardKey.java | 40 +++++ ...EntityWithNonDefaultImmutableShardKey.java | 40 +++++ .../ShardedEntityWithNonDefaultShardKey.java | 40 +++++ ...VersionedEntityWithNonDefaultShardKey.java | 43 +++++ .../core/UpdateOperationsUnitTests.java | 150 ++++++++++++++++++ .../BasicMongoPersistentEntityUnitTests.java | 41 +++++ src/main/asciidoc/index.adoc | 1 + src/main/asciidoc/reference/sharding.adoc | 70 ++++++++ 20 files changed, 1093 insertions(+), 24 deletions(-) create mode 100644 spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/mapping/ShardKey.java create mode 100644 spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/mapping/Sharded.java create mode 100644 spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/mapping/ShardingStrategy.java create mode 100644 spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ShardedEntityWithDefaultShardKey.java create mode 100644 spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ShardedEntityWithNonDefaultImmutableShardKey.java create mode 100644 spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ShardedEntityWithNonDefaultShardKey.java create mode 100644 spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ShardedVersionedEntityWithNonDefaultShardKey.java create mode 100644 spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/UpdateOperationsUnitTests.java create mode 100644 src/main/asciidoc/reference/sharding.adoc diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/MappedDocument.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/MappedDocument.java index 72b351fc4..7df9503ca 100644 --- a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/MappedDocument.java +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/MappedDocument.java @@ -80,7 +80,11 @@ public class MappedDocument { } public Bson getIdFilter() { - return Filters.eq(ID_FIELD, document.get(ID_FIELD)); + return new Document(ID_FIELD, document.get(ID_FIELD)); + } + + public Object get(String key) { + return document.get(key); } public UpdateDefinition updateWithoutId() { diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/MongoTemplate.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/MongoTemplate.java index 96b2d4620..4179c0f67 100644 --- a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/MongoTemplate.java +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/MongoTemplate.java @@ -33,7 +33,6 @@ import org.bson.Document; import org.bson.conversions.Bson; import org.slf4j.Logger; import org.slf4j.LoggerFactory; - import org.springframework.beans.BeansException; import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContextAware; @@ -1480,23 +1479,38 @@ public class MongoTemplate implements MongoOperations, ApplicationContextAware, } return execute(collectionName, collection -> { + MongoAction mongoAction = new MongoAction(writeConcern, MongoActionOperation.SAVE, collectionName, entityClass, dbDoc, null); WriteConcern writeConcernToUse = prepareWriteConcern(mongoAction); MappedDocument mapped = MappedDocument.of(dbDoc); + MongoCollection collectionToUse = writeConcernToUse == null // + ? collection // + : collection.withWriteConcern(writeConcernToUse); + if (!mapped.hasId()) { - if (writeConcernToUse == null) { - collection.insertOne(dbDoc); - } else { - collection.withWriteConcern(writeConcernToUse).insertOne(dbDoc); - } - } else if (writeConcernToUse == null) { - collection.replaceOne(mapped.getIdFilter(), dbDoc, new ReplaceOptions().upsert(true)); + collectionToUse.insertOne(dbDoc); } else { - collection.withWriteConcern(writeConcernToUse).replaceOne(mapped.getIdFilter(), dbDoc, - new ReplaceOptions().upsert(true)); + + MongoPersistentEntity entity = mappingContext.getPersistentEntity(entityClass); + UpdateContext updateContext = queryOperations.replaceSingleContext(mapped, true); + Document replacement = updateContext.getMappedUpdate(entity); + + Document filter = updateContext.getMappedQuery(entity); + + if (updateContext.requiresShardKey(filter, entity)) { + + if (entity.getShardKey().isImmutable()) { + filter = updateContext.applyShardKey(entity, filter, null); + } else { + filter = updateContext.applyShardKey(entity, filter, + collection.find(filter, Document.class).projection(updateContext.getMappedShardKey(entity)).first()); + } + } + + collectionToUse.replaceOne(filter, replacement, new ReplaceOptions().upsert(true)); } return mapped.getId(); }); @@ -1615,8 +1629,20 @@ public class MongoTemplate implements MongoOperations, ApplicationContextAware, if (!UpdateMapper.isUpdateObject(updateObj)) { + Document filter = new Document(queryObj); + + if (updateContext.requiresShardKey(filter, entity)) { + + if (entity.getShardKey().isImmutable()) { + filter = updateContext.applyShardKey(entity, filter, null); + } else { + filter = updateContext.applyShardKey(entity, filter, + collection.find(filter, Document.class).projection(updateContext.getMappedShardKey(entity)).first()); + } + } + ReplaceOptions replaceOptions = updateContext.getReplaceOptions(entityClass); - return collection.replaceOne(queryObj, updateObj, replaceOptions); + return collection.replaceOne(filter, updateObj, replaceOptions); } else { return multi ? collection.updateMany(queryObj, updateObj, opts) : collection.updateOne(queryObj, updateObj, opts); diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/QueryOperations.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/QueryOperations.java index 23c88f560..f6852c333 100644 --- a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/QueryOperations.java +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/QueryOperations.java @@ -16,7 +16,10 @@ package org.springframework.data.mongodb.core; import java.util.List; +import java.util.Map; import java.util.Optional; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; import java.util.function.Consumer; import java.util.function.Function; import java.util.stream.Collectors; @@ -154,6 +157,15 @@ class QueryOperations { return new UpdateContext(updateDefinition, query, false, upsert); } + /** + * @param replacement the {@link MappedDocument mapped replacement} document. + * @param upsert use {@literal true} to insert diff when no existing document found. + * @return new instance of {@link UpdateContext}. + */ + UpdateContext replaceSingleContext(MappedDocument replacement, boolean upsert) { + return new UpdateContext(replacement, upsert); + } + /** * Create a new {@link DeleteContext} instance removing all matching documents. * @@ -253,7 +265,6 @@ class QueryOperations { */ Document getMappedSort(@Nullable MongoPersistentEntity entity) { return queryMapper.getMappedSort(query.getSortObject(), entity); - } /** @@ -353,7 +364,6 @@ class QueryOperations { if (ClassUtils.isAssignable(requestedTargetType, propertyType)) { conversionTargetType = propertyType; } - } catch (PropertyReferenceException e) { // just don't care about it as we default to Object.class anyway. } @@ -491,7 +501,9 @@ class QueryOperations { private final boolean multi; private final boolean upsert; - private final UpdateDefinition update; + private final @Nullable UpdateDefinition update; + private final @Nullable MappedDocument mappedDocument; + private final Map, Document> mappedShardKey = new ConcurrentHashMap<>(1); /** * Create a new {@link UpdateContext} instance. @@ -520,6 +532,16 @@ class QueryOperations { this.multi = multi; this.upsert = upsert; this.update = update; + this.mappedDocument = null; + } + + UpdateContext(MappedDocument update, boolean upsert) { + + super(new BasicQuery(new Document(BsonUtils.asMap(update.getIdFilter())))); + this.multi = false; + this.upsert = upsert; + this.mappedDocument = update; + this.update = null; } /** @@ -544,7 +566,7 @@ class QueryOperations { UpdateOptions options = new UpdateOptions(); options.upsert(upsert); - if (update.hasArrayFilters()) { + if (update != null && update.hasArrayFilters()) { options .arrayFilters(update.getArrayFilters().stream().map(ArrayFilter::asDocument).collect(Collectors.toList())); } @@ -602,6 +624,45 @@ class QueryOperations { return mappedQuery; } + Document applyShardKey(@Nullable MongoPersistentEntity domainType, Document filter, + @Nullable Document existing) { + + Document shardKeySource = existing != null ? existing + : mappedDocument != null ? mappedDocument.getDocument() : getMappedUpdate(domainType); + + Document filterWithShardKey = new Document(filter); + for (String key : getMappedShardKeyFields(domainType)) { + if (!filterWithShardKey.containsKey(key)) { + filterWithShardKey.append(key, shardKeySource.get(key)); + } + } + + return filterWithShardKey; + } + + boolean requiresShardKey(Document filter, @Nullable MongoPersistentEntity domainType) { + + if (multi || domainType == null || !domainType.isSharded() || domainType.idPropertyIsShardKey()) { + return false; + } + + if (filter.keySet().containsAll(getMappedShardKeyFields(domainType))) { + return false; + } + + return true; + } + + Set getMappedShardKeyFields(@Nullable MongoPersistentEntity entity) { + return getMappedShardKey(entity).keySet(); + } + + Document getMappedShardKey(@Nullable MongoPersistentEntity entity) { + + return mappedShardKey.computeIfAbsent(entity.getType(), + key -> queryMapper.getMappedFields(entity.getShardKey().getDocument(), entity)); + } + /** * Get the already mapped aggregation pipeline to use with an {@link #isAggregationUpdate()}. * @@ -625,8 +686,11 @@ class QueryOperations { */ Document getMappedUpdate(@Nullable MongoPersistentEntity entity) { - return update instanceof MappedUpdate ? update.getUpdateObject() - : updateMapper.getMappedObject(update.getUpdateObject(), entity); + if (update != null) { + return update instanceof MappedUpdate ? update.getUpdateObject() + : updateMapper.getMappedObject(update.getUpdateObject(), entity); + } + return mappedDocument.getDocument(); } /** diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoTemplate.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoTemplate.java index 63413b77b..b2324f94f 100644 --- a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoTemplate.java +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/ReactiveMongoTemplate.java @@ -39,7 +39,6 @@ import org.reactivestreams.Publisher; import org.reactivestreams.Subscriber; import org.slf4j.Logger; import org.slf4j.LoggerFactory; - import org.springframework.beans.BeansException; import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContextAware; @@ -1638,9 +1637,31 @@ public class ReactiveMongoTemplate implements ReactiveMongoOperations, Applicati ? collection // : collection.withWriteConcern(writeConcernToUse); - Publisher publisher = !mapped.hasId() // - ? collectionToUse.insertOne(document) // - : collectionToUse.replaceOne(mapped.getIdFilter(), document, new ReplaceOptions().upsert(true)); + Publisher publisher = null; + if (!mapped.hasId()) { + publisher = collectionToUse.insertOne(document); + } else { + + MongoPersistentEntity entity = mappingContext.getPersistentEntity(entityClass); + UpdateContext updateContext = queryOperations.replaceSingleContext(mapped, true); + Document filter = updateContext.getMappedQuery(entity); + Document replacement = updateContext.getMappedUpdate(entity); + + Mono theFilter = Mono.just(filter); + + if(updateContext.requiresShardKey(filter, entity)) { + if (entity.getShardKey().isImmutable()) { + theFilter = Mono.just(updateContext.applyShardKey(entity, filter, null)); + } else { + theFilter = Mono.from( + collection.find(filter, Document.class).projection(updateContext.getMappedShardKey(entity)).first()) + .defaultIfEmpty(replacement).map(it -> updateContext.applyShardKey(entity, filter, it)); + } + } + + publisher = theFilter.flatMap( + it -> Mono.from(collectionToUse.replaceOne(it, replacement, updateContext.getReplaceOptions(entityClass)))); + } return Mono.from(publisher).map(o -> mapped.getId()); }); @@ -1778,8 +1799,21 @@ public class ReactiveMongoTemplate implements ReactiveMongoOperations, Applicati if (!UpdateMapper.isUpdateObject(updateObj)) { + Document filter = new Document(queryObj); + Mono theFilter = Mono.just(filter); + + if(updateContext.requiresShardKey(filter, entity)) { + if (entity.getShardKey().isImmutable()) { + theFilter = Mono.just(updateContext.applyShardKey(entity, filter, null)); + } else { + theFilter = Mono.from( + collection.find(filter, Document.class).projection(updateContext.getMappedShardKey(entity)).first()) + .defaultIfEmpty(updateObj).map(it -> updateContext.applyShardKey(entity, filter, it)); + } + } + ReplaceOptions replaceOptions = updateContext.getReplaceOptions(entityClass); - return collectionToUse.replaceOne(queryObj, updateObj, replaceOptions); + return theFilter.flatMap(it -> Mono.from(collectionToUse.replaceOne(it, updateObj, replaceOptions))); } return multi ? collectionToUse.updateMany(queryObj, updateObj, updateOptions) diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/mapping/BasicMongoPersistentEntity.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/mapping/BasicMongoPersistentEntity.java index f24ad0d78..cbc67b13a 100644 --- a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/mapping/BasicMongoPersistentEntity.java +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/mapping/BasicMongoPersistentEntity.java @@ -37,6 +37,7 @@ import org.springframework.expression.spel.standard.SpelExpressionParser; import org.springframework.lang.Nullable; import org.springframework.util.Assert; import org.springframework.util.ClassUtils; +import org.springframework.util.ObjectUtils; import org.springframework.util.StringUtils; /** @@ -63,6 +64,8 @@ public class BasicMongoPersistentEntity extends BasicPersistentEntity extends BasicPersistentEntity extends BasicPersistentEntity extends BasicPersistentEntity entity) { + + if (!entity.isAnnotationPresent(Sharded.class)) { + return ShardKey.none(); + } + + Sharded sharded = entity.getRequiredAnnotation(Sharded.class); + + String[] keyProperties = sharded.shardKey(); + if (ObjectUtils.isEmpty(keyProperties)) { + keyProperties = new String[] { "_id" }; + } + + ShardKey shardKey = ShardingStrategy.HASH.equals(sharded.shardingStrategy()) ? ShardKey.hash(keyProperties) + : ShardKey.range(keyProperties); + + return sharded.immutableKey() ? ShardKey.immutable(shardKey) : shardKey; + } + /** * Handler to collect {@link MongoPersistentProperty} instances and check that each of them is mapped to a distinct * field name. diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/mapping/MongoPersistentEntity.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/mapping/MongoPersistentEntity.java index 8d6aaa150..1cef658d9 100644 --- a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/mapping/MongoPersistentEntity.java +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/mapping/MongoPersistentEntity.java @@ -77,4 +77,38 @@ public interface MongoPersistentEntity extends PersistentEntityShard + * Key used to distribute documents across a sharded MongoDB cluster. + *

+ * {@link ShardKey#isImmutable() Immutable} shard keys indicate a fixed value that is not updated (see + * MongoDB + * Reference: Change a Document’s Shard Key Value), which allows to skip server round trips in cases where a + * potential shard key change might have occurred. + * + * @author Christoph Strobl + * @since 3.0 + */ +public class ShardKey { + + private static final ShardKey NONE = new ShardKey(Collections.emptyList(), null, true); + + private final List propertyNames; + private final @Nullable ShardingStrategy shardingStrategy; + private final boolean immutable; + + private ShardKey(List propertyNames, @Nullable ShardingStrategy shardingStrategy, boolean immutable) { + + this.propertyNames = propertyNames; + this.shardingStrategy = shardingStrategy; + this.immutable = immutable; + } + + /** + * @return the number of properties used to form the shard key. + */ + public int size() { + return propertyNames.size(); + } + + /** + * @return the unmodifiable collection of property names forming the shard key. + */ + public Collection getPropertyNames() { + return propertyNames; + } + + /** + * @return {@literal true} if the shard key of an document does not change. + * @see MongoDB + * Reference: Change a Document’s Shard Key Value + */ + public boolean isImmutable() { + return immutable; + } + + /** + * Get the unmapped MongoDB representation of the {@link ShardKey}. + * + * @return never {@literal null}. + */ + public Document getDocument() { + + Document doc = new Document(); + for (String field : propertyNames) { + doc.append(field, shardingValue()); + } + return doc; + } + + private Object shardingValue() { + return ObjectUtils.nullSafeEquals(ShardingStrategy.HASH, shardingStrategy) ? "hash" : 1; + } + + /** + * {@link ShardKey} indicating no shard key has been defined. + * + * @return {@link #NONE} + */ + public static ShardKey none() { + return NONE; + } + + /** + * Create a new {@link ShardingStrategy#RANGE} shard key. + * + * @param propertyNames must not be {@literal null}. + * @return new instance of {@link ShardKey}. + */ + public static ShardKey range(String... propertyNames) { + return new ShardKey(Arrays.asList(propertyNames), ShardingStrategy.RANGE, false); + } + + /** + * Create a new {@link ShardingStrategy#RANGE} shard key. + * + * @param propertyNames must not be {@literal null}. + * @return new instance of {@link ShardKey}. + */ + public static ShardKey hash(String... propertyNames) { + return new ShardKey(Arrays.asList(propertyNames), ShardingStrategy.HASH, false); + } + + /** + * Turn the given {@link ShardKey} into an {@link #isImmutable() immutable} one. + * + * @param shardKey must not be {@literal null}. + * @return new instance of {@link ShardKey} if the given shard key is not already immutable. + */ + public static ShardKey immutable(ShardKey shardKey) { + + if (shardKey.isImmutable()) { + return shardKey; + } + + return new ShardKey(shardKey.propertyNames, shardKey.shardingStrategy, true); + } +} diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/mapping/Sharded.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/mapping/Sharded.java new file mode 100644 index 000000000..646edfd6a --- /dev/null +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/mapping/Sharded.java @@ -0,0 +1,92 @@ +/* + * Copyright 2020 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 + * + * https://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.mongodb.core.mapping; + +import java.lang.annotation.ElementType; +import java.lang.annotation.Inherited; +import java.lang.annotation.Retention; +import java.lang.annotation.RetentionPolicy; +import java.lang.annotation.Target; + +import org.springframework.core.annotation.AliasFor; +import org.springframework.data.annotation.Persistent; + +/** + * The {@link Sharded} annotation provides meta information about the actual distribution of data across multiple + * machines. The {@link #shardKey()} is used to distribute documents across shards.
+ * Please visit the MongoDB Documentation for more information + * about requirements and limitations of sharding.
+ * Spring Data will automatically add the shard key to filter queries used for + * {@link com.mongodb.client.MongoCollection#replaceOne(org.bson.conversions.Bson, Object)} operations triggered by + * {@code save} operations on {@link org.springframework.data.mongodb.core.MongoOperations} and + * {@link org.springframework.data.mongodb.core.ReactiveMongoOperations} as well as {@code update/upsert} operation + * replacing/upserting a single existing document as long as the given + * {@link org.springframework.data.mongodb.core.query.UpdateDefinition} holds a full copy of the entity.
+ * All other operations that require the presence of the {@literal shard key} in the filter query need to provide the + * information via the {@link org.springframework.data.mongodb.core.query.Query} parameter when invoking the method. + * + * @author Christoph Strobl + * @since 3.0 + */ +@Persistent +@Inherited +@Retention(RetentionPolicy.RUNTIME) +@Target({ ElementType.TYPE }) +public @interface Sharded { + + /** + * Alias for {@link #shardKey()}. + * + * @return {@literal _id} by default. + * @see #shardKey() + */ + @AliasFor("shardKey") + String[] value() default {}; + + /** + * The shard key determines the distribution of the collection’s documents among the cluster’s shards. The shard key + * is either a single or multiple indexed properties that exist in every document in the collection.
+ * By default the {@literal id} property is used for sharding.
+ * NOTE Required indexes will not be created automatically. Use + * {@link org.springframework.data.mongodb.core.index.Indexed} or + * {@link org.springframework.data.mongodb.core.index.CompoundIndex} along with enabled + * {@link org.springframework.data.mongodb.config.MongoConfigurationSupport#autoIndexCreation() auto index creation} + * or set up them up via + * {@link org.springframework.data.mongodb.core.index.IndexOperations#ensureIndex(org.springframework.data.mongodb.core.index.IndexDefinition)}. + * + * @return an empty key by default. Which indicates to use the entities {@literal id} property. + */ + @AliasFor("value") + String[] shardKey() default {}; + + /** + * The sharding strategy to use for distributing data across sharded clusters. + * + * @return {@link ShardingStrategy#RANGE} by default + */ + ShardingStrategy shardingStrategy() default ShardingStrategy.RANGE; + + /** + * As of MongoDB 4.2 it is possible to change the shard key using update. Using immutable shard keys avoids server + * round trips to obtain an entities actual shard key from the database. + * + * @return {@literal false} by default; + * @see MongoDB + * Reference: Change a Document’s Shard Key Value + */ + boolean immutableKey() default false; + +} diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/mapping/ShardingStrategy.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/mapping/ShardingStrategy.java new file mode 100644 index 000000000..922342c49 --- /dev/null +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/core/mapping/ShardingStrategy.java @@ -0,0 +1,35 @@ +/* + * Copyright 2020 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 + * + * https://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.mongodb.core.mapping; + +/** + * @author Christoph Strobl + * @since 3.0 + */ +public enum ShardingStrategy { + + /** + * Ranged sharding involves dividing data into ranges based on the shard key values. Each chunk is then assigned a + * range based on the shard key values. + */ + RANGE, + + /** + * Hashed Sharding involves computing a hash of the shard key field’s value. Each chunk is then assigned a range based + * on the hashed shard key values. + */ + HASH +} diff --git a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/util/BsonUtils.java b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/util/BsonUtils.java index 6ecd474d1..7cd569ed4 100644 --- a/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/util/BsonUtils.java +++ b/spring-data-mongodb/src/main/java/org/springframework/data/mongodb/util/BsonUtils.java @@ -54,13 +54,15 @@ public class BsonUtils { } public static Map asMap(Bson bson) { + if (bson instanceof Document) { return (Document) bson; } if (bson instanceof BasicDBObject) { return ((BasicDBObject) bson); } - throw new IllegalArgumentException("o_O what's that? Cannot read values from " + bson.getClass()); + + return (Map) bson.toBsonDocument(Document.class, MongoClientSettings.getDefaultCodecRegistry()); } public static void addToMap(Bson bson, String key, @Nullable Object value) { diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/MongoTemplateUnitTests.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/MongoTemplateUnitTests.java index c86eec589..7f93fdc8e 100644 --- a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/MongoTemplateUnitTests.java +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/MongoTemplateUnitTests.java @@ -1829,6 +1829,95 @@ public class MongoTemplateUnitTests extends MongoOperationsUnitTests { assertThat(captor.getValue()).isEqualTo(Collections.singletonList(Document.parse("{ $unset : \"firstname\" }"))); } + @Test // DATAMONGO-2341 + void saveShouldAppendNonDefaultShardKeyIfNotPresentInFilter() { + + template.save(new ShardedEntityWithNonDefaultShardKey("id-1", "AT", 4230)); + + ArgumentCaptor filter = ArgumentCaptor.forClass(Bson.class); + verify(collection).replaceOne(filter.capture(), any(), any()); + + assertThat(filter.getValue()).isEqualTo(new Document("_id", "id-1").append("country", "AT").append("userid", 4230)); + } + + @Test // DATAMONGO-2341 + void saveShouldAppendNonDefaultShardKeyToVersionedEntityIfNotPresentInFilter() { + + when(collection.replaceOne(any(), any(), any(ReplaceOptions.class))).thenReturn(UpdateResult.acknowledged(1, 1L, null)); + + template.save(new ShardedVersionedEntityWithNonDefaultShardKey("id-1", 1L, "AT", 4230)); + + ArgumentCaptor filter = ArgumentCaptor.forClass(Bson.class); + verify(collection).replaceOne(filter.capture(), any(), any()); + + assertThat(filter.getValue()).isEqualTo(new Document("_id", "id-1").append("version", 1L).append("country", "AT").append("userid", 4230)); + } + + @Test // DATAMONGO-2341 + void saveShouldAppendNonDefaultShardKeyFromExistingDocumentIfNotPresentInFilter() { + + when(findIterable.first()).thenReturn(new Document("_id", "id-1").append("country", "US").append("userid", 4230)); + + template.save(new ShardedEntityWithNonDefaultShardKey("id-1", "AT", 4230)); + + ArgumentCaptor filter = ArgumentCaptor.forClass(Bson.class); + ArgumentCaptor replacement = ArgumentCaptor.forClass(Document.class); + + verify(collection).replaceOne(filter.capture(), replacement.capture(), any()); + + assertThat(filter.getValue()).isEqualTo(new Document("_id", "id-1").append("country", "US").append("userid", 4230)); + assertThat(replacement.getValue()).containsEntry("country", "AT").containsEntry("userid", 4230); + } + + @Test // DATAMONGO-2341 + void saveShouldAppendNonDefaultShardKeyFromGivenDocumentIfShardKeyIsImmutable() { + + template.save(new ShardedEntityWithNonDefaultImmutableShardKey("id-1", "AT", 4230)); + + ArgumentCaptor filter = ArgumentCaptor.forClass(Bson.class); + ArgumentCaptor replacement = ArgumentCaptor.forClass(Document.class); + + verify(collection).replaceOne(filter.capture(), replacement.capture(), any()); + + assertThat(filter.getValue()).isEqualTo(new Document("_id", "id-1").append("country", "AT").append("userid", 4230)); + assertThat(replacement.getValue()).containsEntry("country", "AT").containsEntry("userid", 4230); + + verifyNoInteractions(findIterable); + } + + @Test // DATAMONGO-2341 + void saveShouldAppendDefaultShardKeyIfNotPresentInFilter() { + + template.save(new ShardedEntityWithDefaultShardKey("id-1", "AT", 4230)); + + ArgumentCaptor filter = ArgumentCaptor.forClass(Bson.class); + verify(collection).replaceOne(filter.capture(), any(), any()); + + assertThat(filter.getValue()).isEqualTo(new Document("_id", "id-1")); + verify(findIterable, never()).first(); + } + + @Test // DATAMONGO-2341 + void saveShouldProjectOnShardKeyWhenLoadingExistingDocument() { + + when(findIterable.first()).thenReturn(new Document("_id", "id-1").append("country", "US").append("userid", 4230)); + + template.save(new ShardedEntityWithNonDefaultShardKey("id-1", "AT", 4230)); + + verify(findIterable).projection(new Document("country", 1).append("userid", 1)); + } + + @Test // DATAMONGO-2341 + void saveVersionedShouldProjectOnShardKeyWhenLoadingExistingDocument() { + + when(collection.replaceOne(any(), any(), any(ReplaceOptions.class))).thenReturn(UpdateResult.acknowledged(1, 1L, null)); + when(findIterable.first()).thenReturn(new Document("_id", "id-1").append("country", "US").append("userid", 4230)); + + template.save(new ShardedVersionedEntityWithNonDefaultShardKey("id-1", 1L, "AT", 4230)); + + verify(findIterable).projection(new Document("country", 1).append("userid", 1)); + } + class AutogenerateableId { @Id BigInteger id; diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ReactiveMongoTemplateUnitTests.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ReactiveMongoTemplateUnitTests.java index f48d81bb6..240fdb6c3 100644 --- a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ReactiveMongoTemplateUnitTests.java +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ReactiveMongoTemplateUnitTests.java @@ -18,6 +18,7 @@ package org.springframework.data.mongodb.core; import static org.assertj.core.api.Assertions.*; import static org.mockito.Mockito.*; import static org.springframework.data.mongodb.core.aggregation.Aggregation.*; +import static org.springframework.data.mongodb.test.util.Assertions.assertThat; import lombok.Data; import reactor.core.publisher.Mono; @@ -999,6 +1000,101 @@ public class ReactiveMongoTemplateUnitTests { assertThat(captor.getValue()).isEqualTo(Collections.singletonList(Document.parse("{ $unset : \"firstname\" }"))); } + @Test // DATAMONGO-2341 + void saveShouldAppendNonDefaultShardKeyIfNotPresentInFilter() { + + when(findPublisher.first()).thenReturn(Mono.empty()); + + template.save(new ShardedEntityWithNonDefaultShardKey("id-1", "AT", 4230)).subscribe(); + + ArgumentCaptor filter = ArgumentCaptor.forClass(Bson.class); + verify(collection).replaceOne(filter.capture(), any(), any()); + + assertThat(filter.getValue()).isEqualTo(new Document("_id", "id-1").append("country", "AT").append("userid", 4230)); + } + + @Test // DATAMONGO-2341 + void saveShouldAppendNonDefaultShardKeyFromGivenDocumentIfShardKeyIsImmutable() { + + template.save(new ShardedEntityWithNonDefaultImmutableShardKey("id-1", "AT", 4230)).subscribe(); + + ArgumentCaptor filter = ArgumentCaptor.forClass(Bson.class); + ArgumentCaptor replacement = ArgumentCaptor.forClass(Document.class); + + verify(collection).replaceOne(filter.capture(), replacement.capture(), any()); + + assertThat(filter.getValue()).isEqualTo(new Document("_id", "id-1").append("country", "AT").append("userid", 4230)); + assertThat(replacement.getValue()).containsEntry("country", "AT").containsEntry("userid", 4230); + + verifyNoInteractions(findPublisher); + } + + @Test // DATAMONGO-2341 + void saveShouldAppendNonDefaultShardKeyToVersionedEntityIfNotPresentInFilter() { + + when(collection.replaceOne(any(Bson.class), any(Document.class), any(ReplaceOptions.class))) + .thenReturn(Mono.just(UpdateResult.acknowledged(1, 1L, null))); + when(findPublisher.first()).thenReturn(Mono.empty()); + + template.save(new ShardedVersionedEntityWithNonDefaultShardKey("id-1", 1L, "AT", 4230)).subscribe(); + + ArgumentCaptor filter = ArgumentCaptor.forClass(Bson.class); + verify(collection).replaceOne(filter.capture(), any(), any()); + + assertThat(filter.getValue()) + .isEqualTo(new Document("_id", "id-1").append("version", 1L).append("country", "AT").append("userid", 4230)); + } + + @Test // DATAMONGO-2341 + void saveShouldAppendNonDefaultShardKeyFromExistingDocumentIfNotPresentInFilter() { + + when(findPublisher.first()) + .thenReturn(Mono.just(new Document("_id", "id-1").append("country", "US").append("userid", 4230))); + + template.save(new ShardedEntityWithNonDefaultShardKey("id-1", "AT", 4230)).subscribe(); + + ArgumentCaptor filter = ArgumentCaptor.forClass(Bson.class); + ArgumentCaptor replacement = ArgumentCaptor.forClass(Document.class); + + verify(collection).replaceOne(filter.capture(), replacement.capture(), any()); + + assertThat(filter.getValue()).isEqualTo(new Document("_id", "id-1").append("country", "US").append("userid", 4230)); + assertThat(replacement.getValue()).containsEntry("country", "AT").containsEntry("userid", 4230); + } + + @Test // DATAMONGO-2341 + void saveShouldAppendDefaultShardKeyIfNotPresentInFilter() { + + template.save(new ShardedEntityWithDefaultShardKey("id-1", "AT", 4230)).subscribe(); + + ArgumentCaptor filter = ArgumentCaptor.forClass(Bson.class); + verify(collection).replaceOne(filter.capture(), any(), any()); + + assertThat(filter.getValue()).isEqualTo(new Document("_id", "id-1")); + } + + @Test // DATAMONGO-2341 + void saveShouldProjectOnShardKeyWhenLoadingExistingDocument() { + + when(findPublisher.first()).thenReturn(Mono.just(new Document("_id", "id-1").append("country", "US").append("userid", 4230))); + + template.save(new ShardedEntityWithNonDefaultShardKey("id-1", "AT", 4230)).subscribe(); + + verify(findPublisher).projection(new Document("country", 1).append("userid", 1)); + } + + @Test // DATAMONGO-2341 + void saveVersionedShouldProjectOnShardKeyWhenLoadingExistingDocument() { + + when(collection.replaceOne(any(Bson.class), any(Document.class), any(ReplaceOptions.class))) + .thenReturn(Mono.just(UpdateResult.acknowledged(1, 1L, null))); + when(findPublisher.first()).thenReturn(Mono.empty()); + + template.save(new ShardedVersionedEntityWithNonDefaultShardKey("id-1", 1L, "AT", 4230)).subscribe(); + + verify(findPublisher).projection(new Document("country", 1).append("userid", 1)); + } + @Data @org.springframework.data.mongodb.core.mapping.Document(collection = "star-wars") static class Person { diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ShardedEntityWithDefaultShardKey.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ShardedEntityWithDefaultShardKey.java new file mode 100644 index 000000000..0f8413e88 --- /dev/null +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ShardedEntityWithDefaultShardKey.java @@ -0,0 +1,40 @@ +/* + * Copyright 2020 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 + * + * https://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.mongodb.core; + +import lombok.AllArgsConstructor; +import lombok.Data; + +import org.springframework.data.annotation.Id; +import org.springframework.data.mongodb.core.mapping.Field; +import org.springframework.data.mongodb.core.mapping.Sharded; + +/** + * @author Christoph Strobl + */ +@Data +@AllArgsConstructor +@Sharded +public class ShardedEntityWithDefaultShardKey { + + private @Id String id; + + private String country; + + @Field("userid") // + private Integer userId; + +} diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ShardedEntityWithNonDefaultImmutableShardKey.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ShardedEntityWithNonDefaultImmutableShardKey.java new file mode 100644 index 000000000..85e97d5be --- /dev/null +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ShardedEntityWithNonDefaultImmutableShardKey.java @@ -0,0 +1,40 @@ +/* + * Copyright 2020 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 + * + * https://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.mongodb.core; + +import lombok.AllArgsConstructor; +import lombok.Data; + +import org.springframework.data.annotation.Id; +import org.springframework.data.mongodb.core.mapping.Field; +import org.springframework.data.mongodb.core.mapping.Sharded; + +/** + * @author Christoph Strobl + */ +@Data +@AllArgsConstructor +@Sharded(shardKey = { "country", "userId" }, immutableKey = true) +public class ShardedEntityWithNonDefaultImmutableShardKey { + + private @Id String id; + + private String country; + + @Field("userid") // + private Integer userId; + +} diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ShardedEntityWithNonDefaultShardKey.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ShardedEntityWithNonDefaultShardKey.java new file mode 100644 index 000000000..2aa424360 --- /dev/null +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ShardedEntityWithNonDefaultShardKey.java @@ -0,0 +1,40 @@ +/* + * Copyright 2020 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 + * + * https://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.mongodb.core; + +import lombok.AllArgsConstructor; +import lombok.Data; + +import org.springframework.data.annotation.Id; +import org.springframework.data.mongodb.core.mapping.Field; +import org.springframework.data.mongodb.core.mapping.Sharded; + +/** + * @author Christoph Strobl + */ +@Data +@AllArgsConstructor +@Sharded(shardKey = { "country", "userId" }) +public class ShardedEntityWithNonDefaultShardKey { + + private @Id String id; + + private String country; + + @Field("userid") // + private Integer userId; + +} diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ShardedVersionedEntityWithNonDefaultShardKey.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ShardedVersionedEntityWithNonDefaultShardKey.java new file mode 100644 index 000000000..e0b07885a --- /dev/null +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/ShardedVersionedEntityWithNonDefaultShardKey.java @@ -0,0 +1,43 @@ +/* + * Copyright 2020 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 + * + * https://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.mongodb.core; + +import lombok.AllArgsConstructor; +import lombok.Data; + +import org.springframework.data.annotation.Id; +import org.springframework.data.annotation.Version; +import org.springframework.data.mongodb.core.mapping.Field; +import org.springframework.data.mongodb.core.mapping.Sharded; + +/** + * @author Christoph Strobl + */ +@Data +@AllArgsConstructor +@Sharded(shardKey = { "country", "userId" }) +public class ShardedVersionedEntityWithNonDefaultShardKey { + + private @Id String id; + + private @Version Long version; + + private String country; + + @Field("userid") // + private Integer userId; + +} diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/UpdateOperationsUnitTests.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/UpdateOperationsUnitTests.java new file mode 100644 index 000000000..8edffec3e --- /dev/null +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/UpdateOperationsUnitTests.java @@ -0,0 +1,150 @@ +/* + * Copyright 2020 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 + * + * https://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.mongodb.core; + +import static org.assertj.core.api.Assertions.*; + +import java.util.Arrays; + +import org.bson.Document; +import org.junit.jupiter.api.Test; +import org.springframework.data.mongodb.CodecRegistryProvider; +import org.springframework.data.mongodb.core.convert.MappingMongoConverter; +import org.springframework.data.mongodb.core.convert.MongoConverter; +import org.springframework.data.mongodb.core.convert.NoOpDbRefResolver; +import org.springframework.data.mongodb.core.convert.QueryMapper; +import org.springframework.data.mongodb.core.convert.UpdateMapper; +import org.springframework.data.mongodb.core.mapping.MongoMappingContext; +import org.springframework.data.mongodb.core.mapping.MongoPersistentEntity; +import org.springframework.lang.NonNull; +import org.springframework.lang.Nullable; + +import com.mongodb.MongoClientSettings; + +/** + * @author Christoph Strobl + */ +class UpdateOperationsUnitTests { + + static final Document SHARD_KEY = new Document("country", "AT").append("userid", "4230"); + static final Document SOURCE_DOC = appendShardKey(new Document("_id", "id-1")); + + MongoMappingContext mappingContext = new MongoMappingContext(); + MongoConverter mongoConverter = new MappingMongoConverter(NoOpDbRefResolver.INSTANCE, mappingContext); + QueryMapper queryMapper = new QueryMapper(mongoConverter); + UpdateMapper updateMapper = new UpdateMapper(mongoConverter); + EntityOperations entityOperations = new EntityOperations(mappingContext); + + ExtendedQueryOperations queryOperations = new ExtendedQueryOperations(queryMapper, updateMapper, entityOperations, + MongoClientSettings::getDefaultCodecRegistry); + + @Test // DATAMONGO-2341 + void appliesShardKeyToFilter() { + + Document sourceFilter = new Document("name", "kaladin"); + assertThat(shardedFilter(sourceFilter, ShardedEntityWithNonDefaultShardKey.class, null)) + .isEqualTo(appendShardKey(sourceFilter)); + } + + @Test + void applyShardKeyDoesNotAlterSourceFilter() { + + Document sourceFilter = new Document("name", "kaladin"); + shardedFilter(sourceFilter, ShardedEntityWithNonDefaultShardKey.class, null); + assertThat(sourceFilter).isEqualTo(new Document("name", "kaladin")); + } + + @Test // DATAMONGO-2341 + void appliesExistingShardKeyToFilter() { + + Document sourceFilter = new Document("name", "kaladin"); + Document existing = new Document("country", "GB").append("userid", "007"); + + assertThat(shardedFilter(sourceFilter, ShardedEntityWithNonDefaultShardKey.class, existing)) + .isEqualTo(new Document(existing).append("name", "kaladin")); + } + + @Test // DATAMONGO-2341 + void recognizesExistingShardKeyInFilter() { + + Document sourceFilter = appendShardKey(new Document("name", "kaladin")); + + assertThat(queryOperations.replaceSingleContextFor(SOURCE_DOC).requiresShardKey(sourceFilter, + entityOf(ShardedEntityWithNonDefaultShardKey.class))).isFalse(); + } + + @Test // DATAMONGO-2341 + void recognizesIdPropertyAsShardKey() { + + Document sourceFilter = new Document("_id", "id-1"); + + assertThat(queryOperations.replaceSingleContextFor(SOURCE_DOC).requiresShardKey(sourceFilter, + entityOf(ShardedEntityWithDefaultShardKey.class))).isFalse(); + } + + @Test // DATAMONGO-2341 + void returnsMappedShardKey() { + + queryOperations.replaceSingleContextFor(SOURCE_DOC) + .getMappedShardKeyFields(entityOf(ShardedEntityWithDefaultShardKey.class)) + .containsAll(Arrays.asList("country", "userid")); + } + + @NonNull + private Document shardedFilter(Document sourceFilter, Class entity, Document existing) { + return queryOperations.replaceSingleContextFor(SOURCE_DOC).applyShardKey(entity, sourceFilter, existing); + } + + private static Document appendShardKey(Document source) { + + Document target = new Document(source); + target.putAll(SHARD_KEY); + return target; + } + + MongoPersistentEntity entityOf(Class type) { + return mappingContext.getPersistentEntity(type); + } + + class ExtendedQueryOperations extends QueryOperations { + + ExtendedQueryOperations(QueryMapper queryMapper, UpdateMapper updateMapper, EntityOperations entityOperations, + CodecRegistryProvider codecRegistryProvider) { + super(queryMapper, updateMapper, entityOperations, codecRegistryProvider); + } + + @NonNull + private ExtendedUpdateContext replaceSingleContextFor(Document source) { + return new ExtendedUpdateContext(MappedDocument.of(source), true); + } + + MongoPersistentEntity entityOf(Class type) { + return mappingContext.getPersistentEntity(type); + } + + class ExtendedUpdateContext extends UpdateContext { + + ExtendedUpdateContext(MappedDocument update, boolean upsert) { + super(update, upsert); + } + + Document applyShardKey(@Nullable Class domainType, Document filter, @Nullable Document existing) { + return applyShardKey(entityOf(domainType), filter, existing); + } + } + + } +} diff --git a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/mapping/BasicMongoPersistentEntityUnitTests.java b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/mapping/BasicMongoPersistentEntityUnitTests.java index 70c3b9767..07100b940 100644 --- a/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/mapping/BasicMongoPersistentEntityUnitTests.java +++ b/spring-data-mongodb/src/test/java/org/springframework/data/mongodb/core/mapping/BasicMongoPersistentEntityUnitTests.java @@ -255,6 +255,38 @@ public class BasicMongoPersistentEntityUnitTests { assertThat(entity.getCollation()).isEqualTo(org.springframework.data.mongodb.core.query.Collation.of("en_US")); } + @Test // DATAMONGO-2341 + public void detectsShardedEntityCorrectly() { + + assertThat(entityOf(WithDefaultShardKey.class).isSharded()).isTrue(); + assertThat(entityOf(Contact.class).isSharded()).isFalse(); + } + + @Test // DATAMONGO-2341 + public void readsDefaultShardKey() { + + assertThat(entityOf(WithDefaultShardKey.class).getShardKey().getDocument()) + .isEqualTo(new org.bson.Document("_id", 1)); + } + + @Test // DATAMONGO-2341 + public void readsSingleShardKey() { + + assertThat(entityOf(WithSingleShardKey.class).getShardKey().getDocument()) + .isEqualTo(new org.bson.Document("country", 1)); + } + + @Test // DATAMONGO-2341 + public void readsMultiShardKey() { + + assertThat(entityOf(WithMultiShardKey.class).getShardKey().getDocument()) + .isEqualTo(new org.bson.Document("country", 1).append("userid", 1)); + } + + static BasicMongoPersistentEntity entityOf(Class type) { + return new BasicMongoPersistentEntity<>(ClassTypeInformation.from(type)); + } + @Document("contacts") class Contact {} @@ -313,6 +345,15 @@ public class BasicMongoPersistentEntityUnitTests { @Document(collation = "{ 'locale' : 'en_US' }") class WithDocumentCollation {} + @Sharded + class WithDefaultShardKey {} + + @Sharded("country") + class WithSingleShardKey {} + + @Sharded({ "country", "userid" }) + class WithMultiShardKey {} + static class SampleExtension implements EvaluationContextExtension { /* diff --git a/src/main/asciidoc/index.adoc b/src/main/asciidoc/index.adoc index 11533441b..12a439042 100644 --- a/src/main/asciidoc/index.adoc +++ b/src/main/asciidoc/index.adoc @@ -30,6 +30,7 @@ include::reference/reactive-mongo-repositories.adoc[leveloffset=+1] include::{spring-data-commons-docs}/auditing.adoc[leveloffset=+1] include::reference/mongo-auditing.adoc[leveloffset=+1] include::reference/mapping.adoc[leveloffset=+1] +include::reference/sharding.adoc[leveloffset=+1] include::reference/kotlin.adoc[leveloffset=+1] include::reference/cross-store.adoc[leveloffset=+1] include::reference/jmx.adoc[leveloffset=+1] diff --git a/src/main/asciidoc/reference/sharding.adoc b/src/main/asciidoc/reference/sharding.adoc new file mode 100644 index 000000000..2ca3713a2 --- /dev/null +++ b/src/main/asciidoc/reference/sharding.adoc @@ -0,0 +1,70 @@ +[[sharding]] += Sharding + +MongoDB supports large data sets via sharding, a method for distributing data across multiple machines. Please refer to the https://docs.mongodb.com/manual/sharding/[MongoDB Documentation] to learn how to set up a sharded cluster, its requirements and limitations. + +Spring Data MongoDB uses the `@Sharded` annotation to identify entities stored in sharded collections as shown below. + +==== +[source, java] +---- +@Document("users") +@Sharded(shardKey = { "country", "userId" }) <1> +public class User { + + @Id + Long id; + + @Field("userid") + String userId; + + String country; +} +---- +<1> The properties of the shard key are mapped to the actual field names. See +==== + +[[sharding.sharded-collections]] +== Sharded Collections + +Spring Data MongoDB does not auto set up sharding for collections nor indexes required for it. The snippet below shows how to do so using the MongoDB client API. + +==== +[source, java] +---- +MongoDatabase adminDB = template.getMongoDbFactory() + .getMongoDatabase("admin"); <1> + +adminDB.runCommand(new Document("enableSharding", "db")); <2> + +Document shardCmd = new Document("shardCollection", "db.users") <3> + .append("key", new Document("country", 1).append("userid", 1)); <4> + +adminDB.runCommand(shardCmd); +---- +<1> Sharding commands need to be run against the _admin_ database. +<2> Enable sharding for a specific database if necessary. +<3> Shard a collection within the database having sharding enabled. +<4> Set the shard key (Range based sharding in this case). +==== + +[[sharding.shard-key]] +== Shard Key Handling + +The shard key consists of a single or multiple properties present in every document within the target collection, and is used to distribute them across shards. + +Adding the `@Sharded` annotation to an entity enables Spring Data MongoDB to do best effort optimisations required for sharded scenarios when using repositories. +This means essentially adding required shard key information, if not already present, to `replaceOne` filter queries when upserting entities. This may require an additional server round trip to determine the actual value of the current shard key. + +TIP: By setting `@Sharded(immutableKey = true)` no attempt will be made to check if an entities shard key changed. + +Please see the https://docs.mongodb.com/manual/reference/method/db.collection.replaceOne/#upsert[MongoDB Documentation] for further details and the list below for which operations are eligible for auto include the shard key. + +* `Reactive/CrudRepository.save(...)` +* `Reactive/CrudRepository.saveAll(...)` +* `Reactive/MongoTemplate.save(...)` + + + + +