如何在Spring Boot Webflux与ReactiveMongoRepository中安全实现租户键多租户?
Spring Boot Webflux + MongoDB 多租户自动实现方案
首先明确:可以用AspectJ实现,但Spring Data MongoDB本身提供了更适配响应式场景的扩展方案,比AspectJ更优雅、更符合Spring生态设计,避免响应式编程中的线程阻塞问题。以下是具体实现步骤:
1. 租户上下文传递(响应式适配)
WebFlux不能用ThreadLocal存储租户ID,需要通过Mono上下文传递:
public class TenantContext { private static final String TENANT_ID_KEY = "tenantId"; // 获取当前上下文的租户ID public static Mono<String> getCurrentTenantId() { return Mono.deferContextual(ctx -> ctx.getOrEmpty(TENANT_ID_KEY) .map(Mono::just) .orElseThrow(() -> new IllegalStateException("租户ID未在上下文初始化")) ); } // 向Mono上下文注入租户ID public static <T> Mono<T> withTenantId(String tenantId, Mono<T> mono) { return mono.contextWrite(ctx -> ctx.put(TENANT_ID_KEY, tenantId)); } }
2. 保存数据自动设置tenantId
通过MongoDB的响应式事件监听实现,无需修改Repository代码:
@Component public class TenantIdBeforeConvertListener implements ReactiveBeforeConvertEventListener<TenantScopedModel> { @Override public <S extends TenantScopedModel> Mono<Void> onBeforeConvert(BeforeConvertEvent<S> event) { return TenantContext.getCurrentTenantId() .doOnNext(tenantId -> event.getSource().setAccountId(tenantId)) .then(); } }
所有继承TenantScopedModel的实体在保存前,会自动从上下文获取租户ID并赋值。
3. 查询自动添加tenantId过滤条件
方案:自定义Repository基类 + 查询拦截器
3.1 自定义Repository实现基类
继承SimpleReactiveMongoRepository,重写基础查询方法并自动追加租户过滤:
@NoRepositoryBean public class TenantScopedRepositoryImpl<T extends TenantScopedModel, ID> extends SimpleReactiveMongoRepository<T, ID> implements TenantScopedRepository<T, ID> { private final ReactiveMongoOperations mongoOperations; private final EntityInformation<T, ID> entityInformation; public TenantScopedRepositoryImpl(EntityInformation<T, ID> entityInformation, ReactiveMongoOperations mongoOperations) { super(entityInformation, mongoOperations); this.mongoOperations = mongoOperations; this.entityInformation = entityInformation; } @Override public Mono<T> findById(ID id) { return TenantContext.getCurrentTenantId() .flatMap(tenantId -> mongoOperations.findById(id, entityInformation.getJavaType()) .filter(entity -> tenantId.equals(entity.getAccountId()))); } @Override public Flux<T> findAll() { return TenantContext.getCurrentTenantId() .flatMapMany(tenantId -> mongoOperations.query(entityInformation.getJavaType()) .matching(Criteria.where("accountId").is(tenantId)) .all()); } }
3.2 配置Spring Data使用自定义基类
@Configuration public class MongoConfig extends AbstractReactiveMongoConfiguration { @Override protected String getDatabaseName() { return "your_database_name"; } @Override public MongoClient reactiveMongoClient() { return MongoClients.create(); } @Override protected RepositoryFactorySupport getRepositoryFactory(ReactiveMongoOperations operations) { ReactiveMongoRepositoryFactory factory = new ReactiveMongoRepositoryFactory(operations); // 指定自定义Repository基类 factory.setRepositoryBaseClass(TenantScopedRepositoryImpl.class); return factory; } }
3.3 拦截自定义查询方法(如findByName)
通过ReactiveMongoQueryExecutionInterceptor自动给所有查询追加租户条件:
@Component public class TenantQueryInterceptor implements ReactiveMongoQueryExecutionInterceptor { @Override public <T> Mono<T> intercept(Query query, Class<T> type, ReactiveMongoQueryExecution execution) { return TenantContext.getCurrentTenantId() .flatMap(tenantId -> { query.addCriteria(Criteria.where("accountId").is(tenantId)); return execution.execute(query, type); }); } }
4. 最终使用方式
只需让业务Repository继承TenantScopedRepository即可,完全复用Spring Boot接口式开发的便捷性:
public interface RoleRepository extends TenantScopedRepository<Role, String> { // 自动追加accountId过滤条件 Mono<Role> findByRoleName(String roleName); }
调用roleRepository.findByRoleName("admin")时,会自动拼接accountId = 当前租户ID的过滤条件。
关于AspectJ的补充说明
如果一定要用AspectJ实现,需要注意响应式类型(Mono/Flux)的处理,但会存在线程阻塞风险(如Flux.filter中调用block()),不推荐在WebFlux场景使用。示例代码仅作参考:
@Aspect @Component public class TenantRepositoryAspect { @Around("execution(* com.your.package.TenantScopedRepository.*(..)) && !execution(* save*(..))") public Object interceptQueries(ProceedingJoinPoint joinPoint) throws Throwable { Object result = joinPoint.proceed(); if (result instanceof Mono) { return ((Mono<?>) result).flatMap(entity -> TenantContext.getCurrentTenantId() .filter(tenantId -> tenantId.equals(((TenantScopedModel) entity).getAccountId())) .thenReturn(entity) ); } // Flux类型处理容易出现阻塞问题,不建议使用 return result; } @Around("execution(* com.your.package.TenantScopedRepository.save*(..))") public Object interceptSaves(ProceedingJoinPoint joinPoint) throws Throwable { Object[] args = joinPoint.getArgs(); return TenantContext.getCurrentTenantId() .flatMap(tenantId -> { for (Object arg : args) { if (arg instanceof TenantScopedModel) { ((TenantScopedModel) arg).setAccountId(tenantId); } } try { return Mono.just(joinPoint.proceed(args)); } catch (Throwable throwable) { return Mono.error(throwable); } }); } }
内容的提问来源于stack exchange,提问作者Arpan Mukherjee
相关产品推荐
相关产品推荐

