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

WebFlux中使用Repository实现多集合聚合的方法咨询

WebFlux + MongoDB: 多集合聚合操作实现方案

Hey there! Let's break down how to perform multi-collection aggregation (for your countries and airports collections) using WebFlux and MongoDB's reactive repositories. The core here is leveraging MongoDB's aggregation framework—specifically the $lookup stage to join collections—paired with WebFlux's non-blocking reactive model.

Step 1: Define Your Entity & DTO Classes

First, make sure your entity classes map correctly to MongoDB collections, and create a DTO to hold the aggregated results (since we're combining data from two collections).

Entity Classes

// 对应countries集合
@Document(collection = "countries")
public class Country {
    @Id
    private String id;
    private String code; // 比如 "CN"、"US",用来关联airports的countryCode
    private String name;

    // 构造函数、getter、setter省略
}

// 对应airports集合
@Document(collection = "airports")
public class Airport {
    @Id
    private String id;
    private String name;
    private String countryCode; // 关联Country的code字段
    private String city;

    // 构造函数、getter、setter省略
}

Aggregation Result DTO

This class will hold the combined data from both collections:

public class CountryWithAirports {
    private String countryId;
    private String countryName;
    private List<Airport> airports;

    // 构造函数、getter、setter省略
}

Option 1: Use @Aggregation Annotation in ReactiveMongoRepository

If your aggregation logic is static, you can define the pipeline directly in your repository interface using the @Aggregation annotation. This is clean and straightforward for fixed queries.

public interface CountryRepository extends ReactiveMongoRepository<Country, String> {

    // 根据国家code查询该国家及其所有机场
    @Aggregation(pipeline = {
        // 可选:先过滤出指定国家(去掉这行就是查询所有国家)
        "{ '$match': { 'code': ?0 } }",
        // 关联airports集合:用country的code匹配airport的countryCode,结果存入airports字段
        "{ '$lookup': { " +
            "'from': 'airports', " +
            "'localField': 'code', " +
            "'foreignField': 'countryCode', " +
            "'as': 'airports' " +
        "} }",
        // 重命名字段并排除不需要的_id
        "{ '$project': { " +
            "'countryId': '$_id', " +
            "'countryName': '$name', " +
            "'airports': 1, " +
            "'_id': 0 " +
        "} }"
    })
    Mono<CountryWithAirports> findCountryWithAirportsByCode(String countryCode);
}

Option 2: Manually Build Aggregation Pipeline (For Dynamic Logic)

If you need dynamic aggregation conditions (like runtime filters or variable stages), use ReactiveMongoTemplate to build the pipeline programmatically. This gives you full control over the aggregation stages.

First, inject ReactiveMongoTemplate:

@Autowired
private ReactiveMongoTemplate mongoTemplate;

Then build and execute the aggregation:

// 查询指定国家及其所有机场
public Mono<CountryWithAirports> getCountryWithAirports(String countryCode) {
    Aggregation aggregation = Aggregation.newAggregation(
        // 过滤条件:匹配指定国家code
        Aggregation.match(Criteria.where("code").is(countryCode)),
        // 关联airports集合
        Aggregation.lookup("airports", "code", "countryCode", "airports"),
        // 投影:重命名字段,排除原生_id
        Aggregation.project()
            .and("_id").as("countryId")
            .and("name").as("countryName")
            .and("airports").include()
            .andExclude("_id")
    );

    // 执行聚合:源集合是countries,结果映射到CountryWithAirports
    return mongoTemplate.aggregate(aggregation, "countries", CountryWithAirports.class)
        .single(); // 返回单个结果(如果查所有就去掉single(),返回Flux)
}

// 查询所有国家及其机场
public Flux<CountryWithAirports> getAllCountriesWithAirports() {
    Aggregation aggregation = Aggregation.newAggregation(
        Aggregation.lookup("airports", "code", "countryCode", "airports"),
        Aggregation.project()
            .and("_id").as("countryId")
            .and("name").as("countryName")
            .and("airports").include()
            .andExclude("_id")
    );

    return mongoTemplate.aggregate(aggregation, "countries", CountryWithAirports.class);
}

Key Notes to Remember

  • Field Type Consistency: Ensure the fields used for joining (code in countries and countryCode in airports) have the same data type (both String, for example)—otherwise $lookup won't match any documents.
  • Add More Stages: You can extend the pipeline with other aggregation stages like $unwind (to flatten airport lists), $group (to aggregate stats), or $sort (to order results) as needed.
  • Reactive Flow: All operations return Mono/Flux—make sure to handle the reactive stream correctly in your service or controller (e.g., using subscribe() or returning it directly in a WebFlux endpoint).

内容的提问来源于stack exchange,提问作者Pheonix

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 08:33:01