如何从Scala其他文件获取异步计算后的mapData数据?
问题
我需要从Calculation.scala文件中获取mapData数据,在Final.scala中的调用方式如下。调用fetch_data()后输出为Future(<not completed>),且result为空。我不清楚如何等待数据全部下载完成后再访问mapData?我是Scala新手,熟悉C++的回调机制,但对Scala的Future、Await、OnComplete的用法不明确,希望得到解决方法。
相关代码
Final.scala
object finalComputation { val calculationInfo = new Calculation() calculationInfo.fetch_data() val result = calculationInfo.getMapData() def main(args: Array[String]): Unit = { // ... 业务逻辑 } }
Calculation.scala
class Calculation { var mapData = Map.empty[String, String] def createMapData(metricItem: ActualMetrics) = { mapData += (metricItem._1 -> metricItem._2) } def getMapData() = { mapData } def fetch_data() = { val totalData: Future[Done] = querApi .getData() .map { data => (data) } } Await.result(totalData, Duration.Inf).runForeach(unit => { createMapData(parse.From(totalData)) }) }
解决方法
1. 修正Calculation.scala的异步逻辑错误
你当前代码存在两处核心问题:Await.result的位置非法(不在方法内),且对异步流的处理逻辑混乱。以下是修正后的版本:
import scala.concurrent.Future import scala.concurrent.ExecutionContext.Implicits.global class Calculation { private var mapData = Map.empty[String, String] def createMapData(metricItem: ActualMetrics): Unit = { mapData += (metricItem._1 -> metricItem._2) } def getMapData(): Map[String, String] = mapData // 假设querApi.getData()返回Future[Seq[ActualMetrics]],如果是Akka Stream Source需另行调整 def fetch_data(): Future[Unit] = { querApi.getData().map { dataList => // 遍历异步返回的数据,填充mapData dataList.foreach(createMapData) } } }
2. 在Final.scala中等待异步任务完成
异步操作不能调用后立刻获取结果,必须等待Future完成,有两种常用方式:
方式一:用Await.result阻塞等待(适合简单同步场景)
import scala.concurrent.Await import scala.concurrent.duration.Duration import scala.concurrent.ExecutionContext.Implicits.global object finalComputation { def main(args: Array[String]): Unit = { val calculationInfo = new Calculation() // 阻塞当前线程,直到fetch_data的异步任务完成 Await.result(calculationInfo.fetch_data(), Duration.Inf) // 此时mapData已填充完毕 val result = calculationInfo.getMapData() // 后续业务逻辑处理 } }
方式二:用onComplete注册回调(非阻塞,类似C++回调机制)
如果不想阻塞线程,可注册回调函数在任务完成后处理数据:
import scala.concurrent.ExecutionContext.Implicits.global import scala.util.{Success, Failure} object finalComputation { def main(args: Array[String]): Unit = { val calculationInfo = new Calculation() calculationInfo.fetch_data().onComplete { case Success(_) => val result = calculationInfo.getMapData() // 处理成功获取到的mapData case Failure(ex) => // 处理异步任务失败的情况,比如打印异常 ex.printStackTrace() } // 控制台程序需防止JVM提前退出,可根据实际情况调整等待时间 Thread.sleep(5000) } }
3. 核心概念快速梳理
Future:代表异步计算的结果,类似C++的std::future,是Scala异步编程的核心载体。Await.result:阻塞当前线程,直到Future完成,适合简单同步场景,但禁止在非阻塞环境(如Akka Actor)中使用。onComplete:注册回调函数,Future完成(成功/失败)时自动触发,是非阻塞的异步处理方式,和C++回调逻辑一致。- ExecutionContext:Scala异步操作的执行上下文,默认可通过
import scala.concurrent.ExecutionContext.Implicits.global引入全局上下文。
内容的提问来源于stack exchange,提问作者Prajwal_7
相关产品推荐
相关产品推荐

