【发布时间】:2019-10-22 18:36:30
【问题描述】:
需要一些帮助来理解以下 Spark 中的行为(使用 Scala 和 Databricks)
我有一些数据帧(如果重要,则从 S3 读取),并将通过以 1000 个(最多)批量发出 HTTP 发布请求来发送该数据。所以我对数据框进行了重新分区,以确保每个分区的记录不超过 1000 条。另外,为每一行创建了一个 json 列(所以我只需要稍后将它们放入一个数组中)
问题在于提出请求。我使用以下代码创建了以下 Serializable 类
import org.apache.spark.sql.{DataFrame, Row}
import org.apache.http.client.methods.HttpPost
import org.apache.http.impl.client.HttpClientBuilder
import org.apache.http.HttpHeaders
import org.apache.http.entity.StringEntity
import org.apache.commons.io.IOUtils
object postObject extends Serializable{
val client = HttpClientBuilder.create().build()
val post = new HttpPost("https://my-cool-api-endpoint")
post.addHeader(HttpHeaders.CONTENT_TYPE,"application/json")
def makeHttpCall(row: Iterator[Row]) = {
val json_str = """{"people": [""" + row.toSeq.map(x => x.getAs[String]("json")).mkString(",") + "]}"
post.setEntity(new StringEntity(json_str))
val response = client.execute(post)
val entity = response.getEntity()
println(Seq(response.getStatusLine.getStatusCode(), response.getStatusLine.getReasonPhrase()))
println(IOUtils.toString(entity.getContent()))
}
}
现在当我尝试以下操作时:
postObject.makeHttpCall(data.head(2).toIterator)
它就像一个魅力。请求通过,屏幕上有一些输出,我的 API 获取了这些数据。
但是当我尝试把它放在 foreachPartition 中时:
data.foreachPartition { x =>
postObject.makeHttpCall(x)
}
什么都没有发生。屏幕上没有输出,我的 API 中没有任何内容。如果我尝试重新运行它,几乎所有阶段都会跳过。我相信,出于任何原因,它只是懒惰地评估我的请求,而不是实际执行它。我不明白为什么,以及如何强迫它。
【问题讨论】:
-
对我来说,您的代码运行良好。您的两个代码 sn-ps 之间的区别在于,第一个在驱动程序上运行,第二个在每个执行程序上运行。您是否在执行程序的类路径中包含了 http 客户端库?
-
@werner 我正在为此使用 Databricks,所以是的,我相信执行程序与驱动程序共享相同的类路径(或者我错了吗?)。另外,如果是这种情况,它不会抛出错误吗?
-
不知道数据块在这里是如何工作的。我使用的是独立集群。我猜你已经检查了执行程序日志?
标签: scala apache-spark serialization httprequest