【发布时间】:2018-05-10 22:23:59
【问题描述】:
我正在尝试从 Spring 应用程序调用长时间运行的方法。此方法通过 JPA 读取数据库并在调用者方法完成并返回时执行其任务。问题是 Spring 应用程序需要这种方法是事务性的,但我就是做不到。我得到的是
Exception in thread "bulk_task_executor_thread1" org.springframework.dao.InvalidDataAccessApiUsageException: You're trying to execute a streaming query method without a surrounding transaction that keeps the connection open so that the Stream can actually be consumed. Make sure the code consuming the stream uses @Transactional or any other way of declaring a (read-only) transaction.
at org.springframework.data.jpa.repository.query.JpaQueryExecution$StreamExecution.doExecute(JpaQueryExecution.java:343)
at org.springframework.data.jpa.repository.query.JpaQueryExecution.execute(JpaQueryExecution.java:87)
at org.springframework.data.jpa.repository.query.AbstractJpaQuery.doExecute(AbstractJpaQuery.java:116)
at org.springframework.data.jpa.repository.query.AbstractJpaQuery.execute(AbstractJpaQuery.java:106)
这是我的代码的 sn-ps:
@Component
public class BulkLoadingService {
private static final Logger LOG = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass());
@Autowired
private LoadingWorkerFactory loadingWorkerFactory;
@Autowired
@Qualifier("bulkTaskExecutor")
private TaskExecutor taskExecutor;
@Autowired
ModelMapper modelMapper;
@Transactional(readOnly = true)
public void loadAllData() {
LOG.info("loadAllData() started.");
LoadingWorker loadingWorker = loadingWorkerFactory.getLoadingWorker();
taskExecutor.execute(loadingWorker);
LOG.info("loadAllData() finished.");
}
}
-
@Configuration
public class BulkFrameworkConfiguration {
public static final int NTHREADS = 10;
@Bean
@Qualifier("bulkTaskExecutor")
public TaskExecutor threadPoolTaskExecutor() {
ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
executor.setCorePoolSize(4);
executor.setMaxPoolSize(NTHREADS);
executor.setThreadNamePrefix("bulk_task_executor_thread");
executor.initialize();
return executor;
}
}
-
import org.modelmapper.ModelMapper;
import org.springframework.messaging.MessageChannel;
@Autowired
private MessageChannel actionBulkOutboundChannel;
@Autowired
private ActionRepository actionRepository;
@Component
public class LoadingWorkerFactory {
@Autowired
ModelMapper modelMapper;
public LoadingWorker getLoadingWorker() {
return new LoadingWorker(actionRepository, channel, modelMapper);
}
}
-
public class LoadingWorker implements Runnable {
private static final Logger LOG = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass());
private ModelMapper modelMapper;
private ActionRepository dataRepository;
private MessageChannel outboundChannel;
public LoadingWorker(ActionRepository repository, MessageChannel channel, ModelMapper modelMapper) {
this.modelMapper = modelMapper;
this.dataRepository = repository;
this.outboundChannel = channel;
}
@Override
@Transactional(readOnly = true)
public void run() {
LOG.info("LoadingWorker.run() started.");
long startTime = System.currentTimeMillis();
long counter = 0;
try(Stream<ActionEntity> entityStream = dataRepository.getAll()) {
counter = entityStream.peek(entity -> {
ActionDocument ad = modelMapper.map(entity, ActionDocument.class);
LOG.debug("About to build a message '{}'", ad);
Message<ActionDocument> message = MessageBuilder.withPayload(ad).build();
try {
outboundChannel.send(message);
} catch (MessagingException me) {
LOG.error("Exception encountered while writing request message to queue: {}", me.getRootCause());
LOG.debug("Exception encountered while writing request message to queue", me);
} catch (Exception e) {
LOG.error("Some exception encountered while writing request message to queue", e);
}
}).count();
}
LOG.info("LoadingWorker.run() finished: {} Documents ({} ms)", counter, System.currentTimeMillis() - startTime);
}
protected Stream<ActionEntity> getEntityStream() {
return dataRepository.getAll();
}
}
-
@Repository
public interface ActionRepository extends JpaRepository<ActionEntity, UUID> {
@Query("SELECT oa FROM OfficeAction oa ORDER By oa.id")
Stream<OfficeAction> getAll();
...
-
@Entity
@Table(
name = "oa_office_action")
public class ActionEntity implements Serializable {
...
简而言之,BulkLoadingService 使用 LoadingWorkerFactory 来创建一个新的 worker(LoadingWorker 实现了 Runnable)并使用 TaskExecutor 来运行这个 worker。 LoadingWorker 尝试访问 ActionRepository 以获取所有内容的 Stream,这就是事情破裂的时候。
如何在 Spring 中创建一个方法(不属于 @Bean 或 @Component 的方法)在另一个线程 @Transactional 中运行? 就我而言,声明式事务支持是不可能的吗?
- 如果我使用 ActionRepository 中的同步 getAll() 方法,该方法使用 List 或 Iterable 而不是 Stream,问题就会消失。但我想加快数据加载速度,所以我使用 Stream。
- 如果我直接从 BulkLoadingService 执行数据加载(无线程),问题就会消失。但我想要异步、多线程执行。
- 在 LoadingWorker 的 run() 方法上添加 @Async 注解没有任何区别。
- 我不希望在 BulkLoadingService 中启动事务并将其扩展到衍生线程中。
附注通过直接使用 TransactionTemplate(而不是尝试使用 @Transactional 注释),我能够让它工作。以下是不同之处:
public class LoadingWorker implements Runnable {
...
private PlatformTransactionManager transactionManager;
...
@Override
public void run() {
LOG.info("LoadingWorker.run() started.");
TransactionTemplate transactionTemplate = new TransactionTemplate(transactionManager);
transactionTemplate.setReadOnly(true);
transactionTemplate.execute(status -> {
LOG.info("Anonymous TransactionCallback started.");
inner();
return null;
});
}
protected void inner() {
long startTime = System.currentTimeMillis();
...
虽然我仍然想知道这是否可以使用注释来完成。
【问题讨论】:
标签: java spring multithreading jpa