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

如何将MongoDB聚合管道存入数据库并支持参数运行与多租户适配?

解决方案:将MongoDB聚合管道存储到数据库并支持动态参数与多租户

1. 设计管道存储集合

在MongoDB中创建report_pipelines集合,用于存储各租户的聚合管道模板,示例结构如下:

{
  "_id": ObjectId("..."),
  "tenantId": "client_001", // 多租户隔离标识
  "reportCode": "monthly_sales", // 报表唯一编码
  "pipelineTemplate": [
    { "$match": { "saleDate": { "$gte": "#{startDate}", "$lte": "#{endDate}" } } },
    { "$group": { "_id": "$productCategory", "totalAmount": { "$sum": "$amount" } } },
    { "$sort": { "totalAmount": -1 } }
  ],
  "description": "月度销售分类报表"
}
  • tenantId:确保每个客户端只能访问自身的报表管道
  • pipelineTemplate:聚合管道模板,用#{参数名}作为占位符,后续替换为实际参数值

2. Spring Data MongoDB实体映射

创建对应实体类映射该集合:

import org.springframework.data.annotation.Id;
import org.springframework.data.mongodb.core.mapping.Document;
import java.util.List;
import java.util.Map;

@Document(collection = "report_pipelines")
public class ReportPipeline {
    @Id
    private String id;
    private String tenantId;
    private String reportCode;
    private List<Map<String, Object>> pipelineTemplate;
    private String description;

    // Getters and Setters
}

3. 管道加载与参数替换

编写服务类,负责加载租户管道、替换动态参数并执行聚合:

import org.springframework.data.mongodb.core.MongoTemplate;
import org.springframework.data.mongodb.core.aggregation.Aggregation;
import org.springframework.stereotype.Service;
import java.util.Date;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.regex.Matcher;
import java.util.regex.Pattern;

@Service
public class ReportService {
    private final MongoTemplate mongoTemplate;
    private final ReportPipelineRepository pipelineRepository;

    public ReportService(MongoTemplate mongoTemplate, ReportPipelineRepository pipelineRepository) {
        this.mongoTemplate = mongoTemplate;
        this.pipelineRepository = pipelineRepository;
    }

    public List<ReportBaseModel> runReport(String tenantId, String reportCode, Date startDate, Date endDate) {
        // 1. 加载对应租户的管道模板(多租户隔离)
        ReportPipeline pipeline = pipelineRepository.findByTenantIdAndReportCode(tenantId, reportCode)
                .orElseThrow(() -> new IllegalArgumentException("报表管道不存在"));

        // 2. 替换模板中的参数占位符
        List<Map<String, Object>> processedPipeline = replaceParameters(pipeline.getPipelineTemplate(), startDate, endDate);

        // 3. 构建并执行聚合
        Aggregation aggregation = Aggregation.newAggregation(processedPipeline);
        return mongoTemplate.aggregate(aggregation, "sales_collection", ReportBaseModel.class).getMappedResults();
    }

    private List<Map<String, Object>> replaceParameters(List<Map<String, Object>> template, Date startDate, Date endDate) {
        Map<String, Object> params = new HashMap<>();
        params.put("startDate", startDate);
        params.put("endDate", endDate);

        Pattern pattern = Pattern.compile("#\\{(\\w+)\\}");
        return template.stream().map(phase -> replacePhaseParams(phase, params, pattern)).toList();
    }

    private Map<String, Object> replacePhaseParams(Map<String, Object> phase, Map<String, Object> params, Pattern pattern) {
        Map<String, Object> processedPhase = new HashMap<>();
        for (Map.Entry<String, Object> entry : phase.entrySet()) {
            Object value = entry.getValue();
            if (value instanceof String) {
                Matcher matcher = pattern.matcher((String) value);
                processedPhase.put(entry.getKey(), matcher.replaceAll(match -> params.get(match.group(1)).toString()));
            } else if (value instanceof Map) {
                processedPhase.put(entry.getKey(), replacePhaseParams((Map<String, Object>) value, params, pattern));
            } else if (value instanceof List) {
                processedPhase.put(entry.getKey(), replaceListParams((List<Object>) value, params, pattern));
            } else {
                processedPhase.put(entry.getKey(), value);
            }
        }
        return processedPhase;
    }

    private List<Object> replaceListParams(List<Object> list, Map<String, Object> params, Pattern pattern) {
        return list.stream().map(item -> {
            if (item instanceof String) {
                Matcher matcher = pattern.matcher((String) item);
                return matcher.replaceAll(match -> params.get(match.group(1)).toString());
            } else if (item instanceof Map) {
                return replacePhaseParams((Map<String, Object>) item, params, pattern);
            } else if (item instanceof List) {
                return replaceListParams((List<Object>) item, params, pattern);
            } else {
                return item;
            }
        }).toList();
    }
}

4. 管道仓库接口

定义Spring Data Repository接口用于查询管道:

import org.springframework.data.mongodb.repository.MongoRepository;
import java.util.Optional;

public interface ReportPipelineRepository extends MongoRepository<ReportPipeline, String> {
    Optional<ReportPipeline> findByTenantIdAndReportCode(String tenantId, String reportCode);
}

5. 关键注意事项

  • 多租户隔离:所有管道查询必须携带tenantId条件,禁止跨租户访问
  • 参数校验:对传入的租户ID、报表编码、日期参数做合法性校验,拦截非法请求
  • 管道合法性验证:加载管道后可通过MongoDB的explain命令验证语法,避免执行错误管道
  • 权限控制:仅允许授权DBA操作report_pipelines集合,防止恶意聚合操作

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 11:45:24