如何将MongoDB聚合管道更新转为带类型检查的Spring Data MongoDB响应式代码?
Spring Data MongoDB 转换 MongoDB findOneAndUpdate 聚合管道实现
问题1:转换为Spring Data MongoDB(非响应式)Kotlin/Java代码
首先定义对应MongoDB user_history集合的实体类:
Java 实体类
@Document(collection = "user_history") public class UserHistory { private String id; private Long userId; private List<HistoryItem> history; private Integer autoIncrement; // 构造器、Getter、Setter } public class HistoryItem { private ObjectId id; private Integer sequence; private Date lastViewedAt; private Date createdAt; // 构造器、Getter、Setter }
Kotlin 实体类
@Document(collection = "user_history") data class UserHistory( val id: String? = null, val userId: Long, val history: List<HistoryItem> = emptyList(), val autoIncrement: Int = 0 ) data class HistoryItem( val id: ObjectId = ObjectId.get(), val sequence: Int, val lastViewedAt: Date = Date(), val createdAt: Date = Date() )
Java 非响应式实现
使用MongoTemplate的findAndModify方法,将MongoDB聚合管道转换为Update对象:
import org.springframework.data.mongodb.core.MongoTemplate; import org.springframework.data.mongodb.core.query.Criteria; import org.springframework.data.mongodb.core.query.Query; import org.springframework.data.mongodb.core.query.Update; import org.springframework.data.mongodb.core.aggregation.Aggregation; import org.springframework.data.mongodb.core.aggregation.ReplaceRootOperation; import org.springframework.data.mongodb.core.aggregation.SetOperation; import org.springframework.data.mongodb.core.query.FindAndModifyOptions; import org.bson.types.ObjectId; import java.util.Date; public class UserHistoryService { private final MongoTemplate mongoTemplate; public UserHistoryService(MongoTemplate mongoTemplate) { this.mongoTemplate = mongoTemplate; } public UserHistory updateUserHistory(Long userId) { // 构建查询条件 Query query = Query.query(Criteria.where("userId").is(userId)); // 构建聚合管道作为更新逻辑 Aggregation aggregation = Aggregation.newAggregation( // $replaceRoot 阶段:合并默认值与现有文档 ReplaceRootOperation.replaceRoot() .withValueOf(Aggregation.mergeObjects( Aggregation.newDocument( "userId", userId, "history", Aggregation.asList(), "autoIncrement", 0 ), "$$ROOT" )), // $set 阶段:更新history和autoIncrement字段 SetOperation.set("history") .toValue(Aggregation.cond( Aggregation.gt(Aggregation.size("$history"), 50), "$history", Aggregation.concatArrays( "$history", Aggregation.asList( Aggregation.newDocument( "_id", new ObjectId(), "sequence", Aggregation.add("$autoIncrement", 1), "lastViewedAt", new Date(), "createdAt", new Date() ) ) ) )), SetOperation.set("autoIncrement") .toValue(Aggregation.add("$autoIncrement", 1)) ); Update update = Update.fromAggregation(aggregation); // 设置操作选项:自动插入不存在的文档,返回更新后的新文档 FindAndModifyOptions options = FindAndModifyOptions.options() .upsert(true) .returnNew(true); return mongoTemplate.findAndModify(query, update, options, UserHistory.class); } }
Kotlin 非响应式实现
import org.springframework.data.mongodb.core.MongoTemplate import org.springframework.data.mongodb.core.query.Criteria import org.springframework.data.mongodb.core.query.Query import org.springframework.data.mongodb.core.query.Update import org.springframework.data.mongodb.core.aggregation.Aggregation import org.springframework.data.mongodb.core.aggregation.ReplaceRootOperation import org.springframework.data.mongodb.core.aggregation.SetOperation import org.springframework.data.mongodb.core.query.FindAndModifyOptions import org.bson.types.ObjectId import java.util.Date class UserHistoryService(private val mongoTemplate: MongoTemplate) { fun updateUserHistory(userId: Long): UserHistory? { val query = Query.query(Criteria.where("userId").`is`(userId)) val aggregation = Aggregation.newAggregation( ReplaceRootOperation.replaceRoot() .withValueOf(Aggregation.mergeObjects( Aggregation.newDocument( "userId" to userId, "history" to emptyList<HistoryItem>(), "autoIncrement" to 0 ), "$$ROOT" )), SetOperation.set("history") .toValue(Aggregation.cond( Aggregation.gt(Aggregation.size("\$history"), 50), "\$history", Aggregation.concatArrays( "\$history", Aggregation.asList( Aggregation.newDocument( "_id" to ObjectId(), "sequence" to Aggregation.add("\$autoIncrement", 1), "lastViewedAt" to Date(), "createdAt" to Date() ) ) ) )), SetOperation.set("autoIncrement") .toValue(Aggregation.add("\$autoIncrement", 1)) ) val update = Update.fromAggregation(aggregation) val options = FindAndModifyOptions.options() .upsert(true) .returnNew(true) return mongoTemplate.findAndModify(query, update, options, UserHistory::class.java) } }
问题2:带类型检查的响应式版本实现
通过实体类字段方法引用替代硬编码字符串,实现类型安全:当@Document实体类的字段名称变更时,会直接触发编译错误。
Java 响应式类型安全实现
import org.springframework.data.mongodb.core.ReactiveMongoTemplate; import org.springframework.data.mongodb.core.query.Criteria; import org.springframework.data.mongodb.core.query.Query; import org.springframework.data.mongodb.core.query.Update; import org.springframework.data.mongodb.core.aggregation.Aggregation; import org.springframework.data.mongodb.core.aggregation.ReplaceRootOperation; import org.springframework.data.mongodb.core.aggregation.SetOperation; import org.springframework.data.mongodb.core.query.FindAndModifyOptions; import org.bson.types.ObjectId; import reactor.core.publisher.Mono; import java.util.Date; public class ReactiveUserHistoryService { private final ReactiveMongoTemplate reactiveMongoTemplate; public ReactiveUserHistoryService(ReactiveMongoTemplate reactiveMongoTemplate) { this.reactiveMongoTemplate = reactiveMongoTemplate; } public Mono<UserHistory> updateUserHistory(Long userId) { // 类型安全查询条件:引用实体类字段,字段变更触发编译错误 Query query = Query.query(Criteria.where(UserHistory::userId).is(userId)); // 获取实体类字段名称,避免硬编码 String userIdField = UserHistory::userId.getName(); String historyField = UserHistory::history.getName(); String autoIncrementField = UserHistory::autoIncrement.getName(); String historyItemIdField = HistoryItem::id.getName(); String sequenceField = HistoryItem::sequence.getName(); String lastViewedAtField = HistoryItem::lastViewedAt.getName(); String createdAtField = HistoryItem::createdAt.getName(); Aggregation aggregation = Aggregation.newAggregation( ReplaceRootOperation.replaceRoot() .withValueOf(Aggregation.mergeObjects( Aggregation.newDocument( userIdField, userId, historyField, Aggregation.asList(), autoIncrementField, 0 ), "$$ROOT" )), SetOperation.set(historyField) .toValue(Aggregation.cond( Aggregation.gt(Aggregation.size("$" + historyField), 50), "$" + historyField, Aggregation.concatArrays( "$" + historyField, Aggregation.asList( Aggregation.newDocument( historyItemIdField, new ObjectId(), sequenceField, Aggregation.add("$" + autoIncrementField, 1), lastViewedAtField, new Date(), createdAtField, new Date() ) ) ) )), SetOperation.set(autoIncrementField) .toValue(Aggregation.add("$" + autoIncrementField, 1)) ); Update update = Update.fromAggregation(aggregation); FindAndModifyOptions options = FindAndModifyOptions.options() .upsert(true) .returnNew(true); return reactiveMongoTemplate.findAndModify(query, update, options, UserHistory.class); } }
Kotlin 响应式类型安全实现
import org.springframework.data.mongodb.core.ReactiveMongoTemplate import org.springframework.data.mongodb.core.query.Criteria import org.springframework.data.mongodb.core.query.Query import org.springframework.data.mongodb.core.query.Update import org.springframework.data.mongodb.core.aggregation.Aggregation import org.springframework.data.mongodb.core.aggregation.ReplaceRootOperation import org.springframework.data.mongodb.core.aggregation.SetOperation import org.springframework.data.mongodb.core.query.FindAndModifyOptions import org.bson.types.ObjectId import reactor.core.publisher.Mono import java.util.Date class ReactiveUserHistoryService(private val reactiveMongoTemplate: ReactiveMongoTemplate) { fun updateUserHistory(userId: Long): Mono<UserHistory> { // 类型安全查询条件 val query = Query.query(Criteria.where(UserHistory::userId).`is`(userId)) // 获取实体类字段名称,避免硬编码 val userIdField = UserHistory::userId.name val historyField = UserHistory::history.name val autoIncrementField = UserHistory::autoIncrement.name val historyItemIdField = HistoryItem::id.name val sequenceField = HistoryItem::sequence.name val lastViewedAtField = HistoryItem::lastViewedAt.name val createdAtField = HistoryItem::createdAt.name val aggregation = Aggregation.newAggregation( ReplaceRootOperation.replaceRoot() .withValueOf(Aggregation.mergeObjects( Aggregation.newDocument( userIdField to userId, historyField to emptyList<HistoryItem>(), autoIncrementField to 0 ), "$$ROOT" )), SetOperation.set(historyField) .toValue(Aggregation.cond( Aggregation.gt(Aggregation.size("\$$historyField"), 50), "\$$historyField", Aggregation.concatArrays( "\$$historyField", Aggregation.asList( Aggregation.newDocument( historyItemIdField to ObjectId(), sequenceField to Aggregation.add("\$$autoIncrementField", 1), lastViewedAtField to Date(), createdAtField to Date() ) ) ) )), SetOperation.set(autoIncrementField) .toValue(Aggregation.add("\$$autoIncrementField", 1)) ) val update = Update.fromAggregation(aggregation) val options = FindAndModifyOptions.options() .upsert(true) .returnNew(true) return reactiveMongoTemplate.findAndModify(query, update, options, UserHistory::class.java) } }
问题3:ReactiveMongoTemplate的等效方法
ReactiveMongoTemplate中没有直接的findOneAndUpdate方法,其等效方法是**findAndModify**。该方法完全对应MongoDB的findOneAndUpdate操作,支持设置upsert、returnNew等选项,返回Mono<T>类型的响应式结果,与非响应式MongoTemplate的findAndModify功能一致。
内容的提问来源于stack exchange,提问作者Sean Hwang
相关产品推荐
相关产品推荐

