【问题标题】:Calling Async REST api from spring batch processor从 spring 批处理器调用 Async REST api
【发布时间】:2018-09-12 20:05:08
【问题描述】:

我写了一个处理列表列表的春季批处理作业。

Reader 返回列表列表。 处理器在每个 ListItem 上工作并返回已处理的 List。 Writer 将 List of List 中的内容写入 DB 和 sftp。

我有一个用例,我从 spring 批处理器调用 Async REST api。 在 ListenableFuture 响应中,我实现了 LitenableFutureCallback 来处理成功和失败,它按预期工作,但在异步调用返回之前,ItemProcessor 不会等待来自异步 api 的回调并将对象(列表)返回给 writer。

我不确定如何实现和处理来自 ItemProcessor 的异步调用。

我确实读过 AsyncItemProcessor 和 AsyncItemWriter,但我不确定在这种情况下是否应该使用它。

我也想过从 AsyncRestTemplate 对 ListenableFuture 响应调用 get(),但根据文档,它会阻塞当前线程,直到它收到响应。

我正在寻求有关如何实现此功能的帮助。代码如下:sn-p:

处理器:

public class MailDocumentProcessor implements ItemProcessor<List<MailingDocsEntity>, List<MailingDocsEntity>> {

... Initialization code

@Override
public List<MailingDocsEntity> process(List<MailingDocsEntity> documentsList) throws Exception {
    logger.info("Entering MailingDocsEntity processor");


    List<MailingDocsEntity> synchronizedList = Collections.synchronizedList(documentsList);


    for (MailingDocsEntity mailingDocsEntity : synchronizedList) {
        System.out.println("Reading Mailing id: " + mailingDocsEntity.getMailingId());

       ..code to get the file

         //If the file is not a pdf convert it
         String fileExtension = readFromSpResponse.getFileExtension();
         String fileName = readFromSpResponse.getFileName();
         byte[] fileBytes = readFromSpResponse.getByteArray();

         try {

             //Do checks to make sure PDF file is being sent
             if (!"pdf".equalsIgnoreCase(fileExtension)) {
                 //Only doc, docx and xlsx conversions are supported

                     ...Building REquest object
                     //make async call to pdf conversion service
            pdfService.convertDocxToPdf(request, mailingDocsEntity);

                 } else {
                     logger.error("The file cannot be converted to a pdf.\n"
                        );

                 }
             }


         } catch (Exception ex){
             logger.error("There has been an exception while processing data", ex);

         }

    }
    return synchronizedList;
}

}

Async PdfConversion 服务类:

@Service
public class PdfService{


   @Autowired
   @Qualifier("MicroServiceAsyncRestTemplate")
   AsyncRestTemplate microServiceAsyncRestTemplate;

   public ConvertDocxToPdfResponse convertDocxToPdf(ConvertDocxToPdfRequest request, MailingDocsEntity mailingDocsEntity){

        ConvertDocxToPdfResponse pdfResponse = new ConvertDocxToPdfResponse();


            try {

                HttpHeaders headers = new HttpHeaders();
                headers.setContentType(MediaType.APPLICATION_JSON);

                HttpEntity<?> entity = new HttpEntity<>(request, headers);



                ListenableFuture<ResponseEntity<ConvertDocxToPdfResponse>> microServiceResponse = microServiceAsyncRestTemplate.postForEntity(batchMailProcessingConfiguration.getPdfUrl(), entity, ConvertDocxToPdfResponse.class);

                ConvertDocxToPdfResponse resultBody = microServiceResponse.get().getBody();
                microServiceResponse.addCallback(new ListenableFutureCallback<ResponseEntity<ConvertDocxToPdfResponse>>()  {

                    @Override
                    public void onSuccess(ResponseEntity<ConvertDocxToPdfResponse> result) {
                        ...code to do stuff on success


                    }

                    @Override
                    public void onFailure(Throwable ex) {
                        pdfResponse.setMessage("Exception while retrieving response");

                    }
                });

            } catch (Exception e) {
                String message = "There has been an error while issuing a pdf generate request to the pdf micro service";
                pdfResponse.setMessage(message);
                logger.error(message, e);
            }


        return pdfResponse;
    }

}

我最初的批处理作业是同步的,我正在转换为异步以加快处理速度。 我确实尝试寻找类似的问题,但找不到足够的信息。 非常感谢任何指针或帮助。

谢谢!!

【问题讨论】:

    标签: java spring spring-batch asyncresttemplate


    【解决方案1】:

    我确实读过 AsyncItemProcessor 和 AsyncItemWriter,但我不确定在这种情况下是否应该使用它。

    是的,AsyncItemProcessorAsyncItemWriter 适合您的用例。 AsyncItemProcessor 将为新线程上的项目执行委托ItemProcessor 的逻辑(您的休息调用)。项目完成后,将结果的Future 传递给要写入的AsynchItemWriter。然后AsynchItemWriter 将打开Future 并写入项目。这些组件的好处是你不必自己处理Futures wrapping、unwrapping 等。

    你可以找到:

    希望这会有所帮助。

    【讨论】:

    • 感谢您的回答。我今天会试试这个。
    • 从我的处理器类异步调用 rest api 不是返回类型。我处理数据库中的一些记录列表,我发送项目中包含的文件以进行 pdf 转换。当它返回时,我更新列表项,它是返回到 Writer 的列表的一部分。我确实实现了 AsyncItemProcessor 和 AsyncItemWriter,但不确定如何将 Future 返回给 writer..
    • 更新:我终于可以从处理器实现 Async Rest api 调用了。我从 AsyncItemProcessor 调用 asynaRest api 并将 Future 返回到我正在处理的 POJO 本身并发送给 writer。在 writer 中,我确实 get() 对 ListenableFuture 对象。我选择 get over 回调方法,因为此时我希望 API 有明确的响应并写回 DB。如果您看到我的方法中的改进范围,请告诉我。
    猜你喜欢
    • 2020-07-26
    • 1970-01-01
    • 2018-11-07
    • 1970-01-01
    • 1970-01-01
    • 2021-10-24
    • 2021-10-25
    • 2021-11-19
    • 1970-01-01
    相关资源
    最近更新 更多