【问题标题】:Making HTTP post requests on Spark usign foreachPartition在 Spark 上使用 foreachPartition 发出 HTTP 发布请求
【发布时间】: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


【解决方案1】:

postObject 有 2 个字段:clientpost,它们必须被序列化。

我不确定client 是否正确序列化。 post 对象可能从多个分区(在同一个工作程序上)发生变异。很多事情都可能在这里出错。

我建议尝试删除 postObject 并将其正文直接内联到 foreachPartition

加法:

尝试自己运行它:

sc.parallelize((1 to 10).toList).foreachPartition(row => {
        val client = HttpClientBuilder.create().build()
        val post = new HttpPost("https://google.com")
        post.addHeader(HttpHeaders.CONTENT_TYPE,"application/json")
        val json_str = """{"people": [""" + row.toSeq.map(x => x.toString).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()))
      })

在本地和集群中运行它。 它成功完成并将 405 错误打印到工作日志。 所以请求肯定会到达服务器。

foreachPartition 不返回任何结果。要调试您的问题,您可以将其更改为 mapPartitions:

val responseCodes = sc.parallelize((1 to 10).toList).mapPartitions(row => {
        val client = HttpClientBuilder.create().build()
        val post = new HttpPost("https://google.com")
        post.addHeader(HttpHeaders.CONTENT_TYPE,"application/json")
        val json_str = """{"people": [""" + row.toSeq.map(x => x.toString).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()))
        Iterator.single(response.getStatusLine.getStatusCode)
      }).collect()

println(responseCodes.mkString(", "))

此代码返回响应代码列表,以便您对其进行分析。 对我来说,它按预期打印405, 405

【讨论】:

  • 嗯,为什么客户端不能正确序列化?该对象扩展了Serializable,并且它没有抛出任何错误(如果我不扩展序列化,那么是的,它会抛出一个错误,即任务不可序列化)。无论如何,尝试将代码直接放入 foreachPartition 中,并得到相同的结果。没有实际发出请求,没有抛出错误,控制台没有输出
  • 我认为 simpadjo 在这里有所作为。当您每个节点使用多个内核时,您的 HttpClient 可能会被多个线程同时使用,并且它可能不是线程安全的。 post 也一样。
  • 看起来我被控制台上缺少反馈所愚弄(println 在 foreachPartition 循环中无法正常工作)。现在我看到一些 422,这就解释了为什么我在另一边看不到任何东西。我会做更多的事情,因为我相信对象方法也可以工作,并在这里发回。谢谢楼主
  • 不客气!对象本身并不坏,我只是想尽可能多地排除问题。但是共享post 对象必须修复。
【解决方案2】:

有一种方法可以做到这一点,而不必找出究竟什么是不可序列化的。如果要保留代码的结构,可以将所有字段设为@transient lazy val。此外,任何具有副作用的调用都应包装在一个块中。例如

val post = {
  val httpPost = new HttpPost("https://my-cool-api-endpoint")
  httpPost.addHeader(HttpHeaders.CONTENT_TYPE,"application/json")
  httpPost
}

这将延迟所有字段的初始化,直到它们被工作人员使用。每个工作人员都有一个对象实例,您将能够调用makeHttpCall 方法。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-02-29
    相关资源
    最近更新 更多