Vin*_*wak 9 scala future apache-spark
我从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的最佳实践是什么.
我认为这个用例看起来很常见,但还没有找到任何答案.
最好的祝福
这个用例看起来很常见
不是真的,因为它根本无法正常工作(可能).由于每个任务都在标准Scala上Iterators运行,因此这些操作将被压缩在一起.这意味着所有操作都将在实践中阻塞.假设您有三个URL ["x","y","z"],您的代码将按以下顺序执行:
Await.result(httpCall("x", 10.seconds))
Await.result(httpCall("y", 10.seconds))
Await.result(httpCall("z", 10.seconds))
Run Code Online (Sandbox Code Playgroud)
您可以轻松地在本地重现相同的行为.如果要异步执行代码,则应使用以下方法显式处理 mapPartitions:
rdd.mapPartitions(iter => {
??? // Submit requests
??? // Wait until all requests completed and return Iterator of results
})
Run Code Online (Sandbox Code Playgroud)
但这相当棘手.无法保证给定分区的所有数据都适合内存,因此您可能也需要一些批处理机制.
所有这些说我无法重现你所描述的问题可能是一些配置问题或httpCall自身的问题.
在旁注上允许单个超时杀死整个任务看起来不是一个好主意.
我找不到一种简单的方法来实现这一目标。但经过多次重试迭代后,这就是我所做的,并且它适用于大量查询。基本上我们用它来对一个巨大的查询执行批处理操作到多个子查询。
// Break down your huge workload into smaller chunks, in this case huge query string is broken
// down to a small set of subqueries
// Here if needed to optimize further down, you can provide an optimal partition when parallelizing
val queries = sqlContext.sparkContext.parallelize[String](subQueryList.toSeq)
// Then map each one those to a Spark Task, in this case its a Future that returns a string
val tasks: RDD[Future[String]] = queries.map(query => {
val task = makeHttpCall(query) // Method returns http call response as a Future[String]
task.recover {
case ex => logger.error("recover: " + ex.printStackTrace()) }
task onFailure {
case t => logger.error("execution failed: " + t.getMessage) }
task
})
// Note:: Http call is still not invoked, you are including this as part of the lineage
// Then in each partition you combine all Futures (means there could be several tasks in each partition) and sequence it
// And Await for the result, in this way you making it to block untill all the future in that sequence is resolved
val contentRdd = tasks.mapPartitions[String] { f: Iterator[Future[String]] =>
val searchFuture: Future[Iterator[String]] = Future sequence f
Await.result(searchFuture, threadWaitTime.seconds)
}
// Note: At this point, you can do any transformations on this rdd and it will be appended to the lineage.
// When you perform any action on that Rdd, then at that point,
// those mapPartition process will be evaluated to find the tasks and the subqueries to perform a full parallel http requests and
// collect those data in a single rdd.
Run Code Online (Sandbox Code Playgroud)
如果您不想对内容执行任何转换,例如解析响应负载等。那么您可以使用foreachPartition而不是mapPartitions立即执行所有这些 http 调用。