【问题标题】:Spring @Transactional support within a thread线程内的 Spring @Transactional 支持
【发布时间】: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


    【解决方案1】:

    无需重构应用程序的整个部分。我想说,由于您已经将 bean 传递给 Runnable 实现,您可以传递 TransactionTemplate 并在与它的事务中执行您的代码。

    您必须使用execute 方法并将实现放入其中。

    【讨论】:

      猜你喜欢
      • 2017-05-06
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-07-03
      • 2011-02-24
      • 1970-01-01
      • 2023-03-31
      相关资源
      最近更新 更多