【问题标题】:Spring dynamic asyncSpring动态异步
【发布时间】:2020-09-20 19:10:50
【问题描述】:

尝试在 Spring 2.0.5 中运行动态数量的异步进程。我确实让它在一个固定的数字上工作,但需要灵活性。

目前的方法是:

在 AsyncFactoryService 类中:

   public void doAsyncBetter(long rowsPerBucket) throws InterruptedException, SQLException, ExecutionException {
    List<Bucket> buckets = sourceData01Service.bucketsByRowsPerBucket(rowsPerBucket);
    AsyncDataLoadService[] asyncDataLoadServiceArray = new AsyncDataLoadService[buckets.size()];
    for(int i = 0; i < buckets.size(); i++) {
      asyncDataLoadServiceArray[i] = new AsyncDataLoadService();
      asyncDataLoadServiceArray[i].loadSourceData01(buckets.get(i));
    }
  }

在类 AsyncDataLoadService 中

  @Autowired SourceData01Service sourceData01Service;
  
  @Async
  public Future<String> loadSourceData01(Bucket bucket) throws InterruptedException, SQLException {
    LOGGER.info(String.format("(loadSourceData01) [Thread id => %d, bucket => %s]",Thread.currentThread().getId(), bucket.toString()));
    List<SourceData01> sourceData01s = null;
    sourceData01s = sourceData01Service.rowsByBucket(bucket);
    Thread.sleep(5000);
    String info = String.format("finish sourceData01s.size() => %d, Thread id => %d", sourceData01s.size(), Thread.currentThread().getId());
    return new AsyncResult<>(info);
  }

这失败了

Caused by: java.lang.NullPointerException
    at demo.service.AsyncDataLoadService.loadSourceData01(AsyncDataLoadService.java:35) ~[classes/:?]
    at demo.service.AsyncFactoryService.doAsyncBetter(AsyncFactoryService.java:43) ~[classes/:?]
    at demo.RunApp.run(RunApp.java:39) ~[classes/:?]

在服务中似乎 sourceData01Service 为空。

当我尝试使用固定数字时

public void doAsync(long rowsPerBucket) throws InterruptedException, SQLException, ExecutionException {
    List<Bucket> buckets = sourceData01Service.bucketsByRowsPerBucket(rowsPerBucket);
    Future<String> process0 = asyncDataLoadService.loadSourceData01(buckets.get(0));
    Future<String> process1 = asyncDataLoadService.loadSourceData01(buckets.get(1));
    Future<String> process2 = asyncDataLoadService.loadSourceData01(buckets.get(2));
    Future<String> process3 = asyncDataLoadService.loadSourceData01(buckets.get(3));
    Future<String> process4 = asyncDataLoadService.loadSourceData01(buckets.get(4));
    while(!(process0.isDone()) && !(process2.isDone()) && !(process3.isDone()) && !(process4.isDone())) {
      Thread.sleep(2000);
    }
    LOGGER.info("Process 0 => " + process0.get());
    LOGGER.info("Process 1 => " + process1.get());
    LOGGER.info("Process 2 => " + process2.get());
    LOGGER.info("Process 3 => " + process3.get());
    LOGGER.info("Process 4 => " + process4.get());
  }

它按预期工作,但在我的情况下,桶的数量取决于表中的行数,所以我不知道会有多少桶。任何想法如何做到这一点?

【问题讨论】:

    标签: spring asynchronous


    【解决方案1】:

    AsyncDataLoadService 中的服务 (sourceData01Service) 注入没有发生,因为应用程序正在创建 AsyncDataLoadService 的实例(在工厂类中)而不是 spring 创建它。

    我们可以修改实现让spring创建需要的实例。

    @Component
    class AsyncDataLoadService {
          @Autowired SourceData01Service sourceData01Service;
          
          @Async
          public Future<String> loadSourceData01(Bucket bucket) throws InterruptedException, SQLException {
            LOGGER.info(String.format("(loadSourceData01) [Thread id => %d, bucket => %s]",Thread.currentThread().getId(), bucket.toString()));
            List<SourceData01> sourceData01s = null;
            sourceData01s = sourceData01Service.rowsByBucket(bucket);
            Thread.sleep(5000);
            String info = String.format("finish sourceData01s.size() => %d, Thread id => %d", sourceData01s.size(), Thread.currentThread().getId());
            return new AsyncResult<>(info);
          }
     }
    

    然后我们可以将工厂类实现更改为 autowire AsyncDataLoadService 并提交要执行的作业(在本例中为 loadSourceData01)。

    由于asyncDataLoadService.loadSourceData01 使用@Async 进行注释,因此对该方法的每次调用都将在单独的线程中执行。默认情况下,spring 使用 SimpleAsyncTaskExecutor 为每次调用启动新线程。如果需要,配置线程池实现。

    @Component
    class AsyncFactoryService {
    
       @Autowired
       AsyncDataLoadService asyncDataLoadService;
    
        public void doAsyncBetter(long rowsPerBucket) throws InterruptedException, SQLException, ExecutionException {
            List<Bucket> buckets = sourceData01Service.bucketsByRowsPerBucket(rowsPerBucket);
            for(int i = 0; i < buckets.size(); i++) {
              asyncDataLoadService.loadSourceData01(buckets.get(i));
            }
        }
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2014-08-24
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-07-21
      • 2013-04-08
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多