我从URL列表中构建一个RDD,然后尝试使用一些异步http调用来获取数据.在进行其他计算之前我需要所有结果.理想情况下,我需要在不同节点上进行http调用以进行扩展考虑.
我做了这样的事情:
//init spark
val sparkContext = new SparkContext(conf)
val datas = Seq[String]("url1", "url2")
//create rdd
val rdd = sparkContext.parallelize[String](datas)
//httpCall return Future[String]
val requests = rdd.map((url: String) => httpCall(url))
//await all results (Future.sequence may be better)
val responses = requests.map(r => Await.result(r, 10.seconds))
//print responses
response.collect().foreach((s: String) => println(s))
//stop spark
sparkContext.stop()
Run Code Online (Sandbox Code Playgroud)
这项工作,但Spark工作永远不会完成!
所以我想知道使用Spark(或Future [RDD])处理Future的最佳实践是什么.
我认为这个用例看起来很常见,但还没有找到任何答案.
最好的祝福