异步流程中Hibernate Filter未生效问题排查求助
问题:Hibernate过滤器在异步线程中未生效,切面已触发但过滤器不工作
我在实体上定义了Hibernate过滤器,通过切面为所有实现TenantableRepository的仓库注入该过滤器。目前遇到的问题是:在CompletableFuture异步方法中执行仓库调用时,过滤器未生效;但在主线程中的仓库调用,过滤器能正常注入。我知道线程不同,但两次调用均触发了切面(有日志打印)。我希望过滤器在API请求和异步流程等所有场景中都能生效。
切面代码
@Aspect @Component @Slf4j public class TenantFilterAspect { @PersistenceContext private EntityManager entityManager; @Before("execution(* com.demo.repository.TenantableRepository+.find*(..))") public void beforeFindOfTenantableRepository() { log.info("Called aspect"); entityManager .unwrap(Session.class) .enableFilter(Tenantable.TENANT_FILTER_NAME) .setParameter(Tenantable.TENANT_COLUMN, TenantContext.getTenantId()); } }
测试控制器代码
import org.springframework.beans.factory.annotation.Autowired; @RestController @Slf4j @RequestMapping(value = "/v1/api/test", consumes = MediaType.APPLICATION_JSON_VALUE, produces = MediaType.APPLICATION_JSON_VALUE) public class TestController { @Autowired MyEntityRepository myEntityRepository; @RequestMapping(value = "/aysnc", method = RequestMethod.GET, consumes = MediaType.ALL_VALUE) public @ResponseBody ResponseEntity<APIResponse> testAsync(HttpServletRequest httpServletRequest) throws InterruptedException { Optional<MyEntity> entity = myEntityRepository.findByFirstName("Firstname"); if(!entity.isEmpty()){ log.info("1. Main: found entity:{}",entity.get()); } CompletableFuture.runAsync(this::callAsyncMethod); return ResponseEntity.status(HttpStatus.OK) .body("Ok"); } public void callAsyncMethod() { Optional<MyEntity> entity = myEntityRepository.findByFirstName("Firstname"); if(!entity.isEmpty()){ log.info("2. Async: found entity:{}",entity.get()); } } }
问题原因
- EntityManager线程绑定限制:你通过
@PersistenceContext注入的EntityManager是线程绑定的代理对象,主线程中它关联的是当前线程的Session,但异步线程中,这个代理会获取全新的Session,你在切面里对主线程Session做的过滤器配置,不会同步到异步线程的新Session上。 - TenantContext的线程隔离:如果
TenantContext是基于ThreadLocal实现的,原生CompletableFuture不会继承主线程的ThreadLocal值,导致异步线程中TenantContext.getTenantId()可能返回null,即使过滤器被启用也无法生效。
修复方案
1. 解决TenantContext的线程传递问题
如果TenantContext用ThreadLocal存储租户ID,有两种处理方式:
- 手动传递租户ID到异步线程:
// 修改控制器中的异步任务提交代码 String tenantId = TenantContext.getTenantId(); CompletableFuture.runAsync(() -> { TenantContext.setTenantId(tenantId); callAsyncMethod(); }); - 改用Spring的
@Async配合配置了inheritableThreadLocals=true的任务执行器,自动继承主线程的ThreadLocal值(推荐)。
2. 切面中操作当前线程的Session
不要直接依赖注入的EntityManager,改为从EntityManagerFactory获取当前线程对应的Session:
@Aspect @Component @Slf4j public class TenantFilterAspect { @Autowired private EntityManagerFactory entityManagerFactory; @Before("execution(* com.demo.repository.TenantableRepository+.find*(..))") public void beforeFindOfTenantableRepository() { log.info("Called aspect"); String tenantId = TenantContext.getTenantId(); if (tenantId == null) { log.warn("Tenant ID is null, skipping filter activation"); return; } // 获取当前线程绑定的EntityManager(事务环境下推荐用这个) EntityManager em = EntityManagerFactoryUtils.getTransactionalEntityManager(entityManagerFactory); if (em == null) { // 非事务场景下创建临时EntityManager em = entityManagerFactory.createEntityManager(); } em.unwrap(Session.class) .enableFilter(Tenantable.TENANT_FILTER_NAME) .setParameter(Tenantable.TENANT_COLUMN, tenantId); } }
3. 改用Spring @Async替代原生CompletableFuture
原生CompletableFuture不会集成Spring上下文,改用Spring的异步机制能更优雅解决线程上下文传递问题:
第一步:配置异步任务执行器
@Configuration @EnableAsync public class AsyncConfig implements AsyncConfigurer { @Override public Executor getAsyncExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(5); executor.setMaxPoolSize(10); executor.setQueueCapacity(25); // 允许继承主线程的ThreadLocal值 executor.setInheritableThreadLocals(true); executor.setThreadNamePrefix("TenantAsync-"); executor.initialize(); return executor; } }
第二步:重构控制器代码
把异步方法移到Spring管理的Bean中,并添加@Async注解:
@RestController @Slf4j @RequestMapping(value = "/v1/api/test", consumes = MediaType.APPLICATION_JSON_VALUE, produces = MediaType.APPLICATION_JSON_VALUE) public class TestController { @Autowired MyEntityRepository myEntityRepository; @Autowired AsyncService asyncService; @RequestMapping(value = "/async", method = RequestMethod.GET, consumes = MediaType.ALL_VALUE) public ResponseEntity<String> testAsync(HttpServletRequest httpServletRequest) { Optional<MyEntity> entity = myEntityRepository.findByFirstName("Firstname"); if(entity.isPresent()){ log.info("1. Main: found entity:{}",entity.get()); } asyncService.callAsyncMethod(); return ResponseEntity.status(HttpStatus.OK).body("Ok"); } } @Service public class AsyncService { @Autowired MyEntityRepository myEntityRepository; @Async public void callAsyncMethod() { Optional<MyEntity> entity = myEntityRepository.findByFirstName("Firstname"); if(entity.isPresent()){ log.info("2. Async: found entity:{}",entity.get()); } } }
内容的提问来源于stack exchange,提问作者user11567166
相关产品推荐
相关产品推荐

