【问题标题】:A better way to handle asynchronous calls when using the AWS Java SDK V2 in Scala?在 Scala 中使用 AWS Java SDK V2 时处理异步调用的更好方法?
【发布时间】:2018-11-26 23:59:43
【问题描述】:

背景

日前,AWS Java SDK 2.1版本正式发布。主要卖点之一是它与以前版本的 SDK 相比如何处理异步调用。

我决定使用 Scala 和新的 SDK 进行一些实验,并且尝试想出一种惯用的方法来处理 SDK 返回的 Futures 有点困难。

问题

有没有一种方法可以更好、更简洁、使用更少的样板代码?

目标

使用 Scala 处理适用于 Java V2 的 AWS 开发工具包,并能够以惯用的方式处理成功和失败。

实验

创建一个Async SNS Client并异步提交消息500:

实验一——使用SDK返回的CompletableFuture

  (0 until 500).map { i =>
    val future = client.publish(PublishRequest.builder().topicArn(arn).message(messageJava + i.toString).build())
    future.whenComplete((response, ex) => {
      val responseOption = Option(response) // Response can be null
      responseOption match {
        case Some(r) => println(r.messageId())
        case None => println(s"There was an error ${ex.getMessage}")
      }
    })
  }.foreach(future => future.join())

在这里,我创建一个独特的请求并发布它。 whenComplete 函数将响应转换为选项,因为该值可以为空。这很丑陋,因为处理成功/失败的方法必须在响应中检查 null。

实验 2 - 在 Scala Future 中获取结果

(0 until 500).map { i =>
    val jf = client.publish(PublishRequest.builder().topicArn(arn).message(messageScala + i.toString).build())
    val sf: Future[PublishResponse] = Future { jf.get }
    sf.onComplete {
      case Success(response) => print(response.messageId)
      case Failure(ex) => println(s"There was an error ${ex.getMessage}")
    }
    sf
  }.foreach(Await.result(_, 5000.millis))

在这里,我在CompletableFuture 上使用.get() 方法,这样我就可以处理Scala Future。

实验 3 - 使用 Scala - Java8 - Compat 库将 CompletableFuture 转换为 Future

(0 until 500).map { i =>
    val f = client.publish(PublishRequest.builder().topicArn(arn).message(messageScala + i.toString).build()).toScala
    f.onComplete {
      case Success(response) =>
      case Failure(exception) => println(exception.getMessage)
    }
    f
  }.foreach(Await.result(_, 5000.millis))

这是迄今为止我最喜欢的实现,除了我需要使用第三方experimental library

观察

  • 总的来说,所有这些实现的性能大致相同,future.join() 比其他实现快一点。
  • 这些函数初始化客户端和发布 500 条消息所用的时间约为 2 秒
  • 此代码的顺序版本需要不到 1 分钟(55 秒)
  • 可以看完整代码here

【问题讨论】:

  • 虽然scala-java8-compat README 仍然声明“API 目前仍处于试验阶段:我们尚不保证与未来版本的源代码或二进制兼容性。”,它是 Scala 作者工作的一部分,因此并不是真正的第三方。

标签: scala amazon-web-services aws-sdk amazon-sns aws-java-sdk


【解决方案1】:

您提到您很高兴将 completablefuture 转换为 scala.future,只是您不喜欢依赖 scala-java8-compat。

在这种情况下,您可以简单地滚动您自己的,并且您只希望 java8 可以 scala:

object CompletableFutureOps {                                                                                                                                        

  implicit class CompletableFutureToScala[T](cf: CompletableFuture[T]) {                                                                                             
    def asScala: Future[T] = {                                                                                                                                       
      val p = Promise[T]()                                                                                                                                           
      cf.whenCompleteAsync{ (result, ex) =>                                                                                                                          
        if (result == null) p failure ex                                                                                                                             
        else                p success result                                                                                                                         
      }                                                                                                                                                              
      p.future                                                                                                                                                       
    }                                                                                                                                                                
  }                                                                                                                                                                  
}

def showByExample: Unit = {
  import CompletableFutureOps._   
  (0 until 500).map { i =>                                                                                                                                                                                                                                                                                     
     val f = CompletableFuture.supplyAsync(() => i).asScala                                                                                                             
     f.onComplete {                                                                                                                                                     
       case Success(response)  => println("Success: " + response)                                                                                                        
       case Failure(exception) => println(exception.getMessage)                                                                                                         
     }                                                                                                                                                                  
     f                                                                                                                                                                  
  }.foreach(Await.result(_, 5000.millis))    
}             

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-09-01
    • 2019-03-19
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多