使用Async HTTP调用进行Spark作业

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的最佳实践是什么.

我认为这个用例看起来很常见,但还没有找到任何答案.

最好的祝福

zer*_*323 9

这个用例看起来很常见

不是真的,因为它根本无法正常工作(可能).由于每个任务都在标准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自身的问题.

在旁注上允许单个超时杀死整个任务看起来不是一个好主意.


rak*_*sja 5

我找不到一种简单的方法来实现这一目标。但经过多次重试迭代后,这就是我所做的,并且它适用于大量查询。基本上我们用它来对一个巨大的查询执行批处理操作到多个子查询。

// 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 调用。

  • 我如何在 python Pyspark 中实现这个解决方案? (2认同)