You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.15 13:12:06