如何将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
相关产品推荐
相关产品推荐

