【问题标题】:Spark foreachpartition connection improvementsSpark foreachpartition 连接改进
【发布时间】:2017-06-20 00:39:52
【问题描述】:

我写了一个火花作业,它在下面的操作中进行

  1. 从 HDFS 文本文件中读取数据。

  2. 执行 distinct() 调用以过滤重复项。

  3. 做一个mapToPair阶段并生成pairRDD

  4. 做一个 reducebykey 调用

  5. 为分组元组做聚合逻辑。

  6. 现在在 #5 上调用 foreach

    在这里

    1. 调用 cassandra db
    2. 创建 aws SNS 和 SQS 客户端连接
    3. 做一些 json 记录格式化。
    4. 将记录发布到 SNS/SQS

当我运行此作业时,它会创建三个火花阶段

第一阶段 - 大约需要 45 秒。执行不同的 第二阶段 - mapToPair 和 reducebykey = 需要 1.5 分钟

第三阶段 = 需要 19 分钟

我做了什么

  1. 我关闭了 cassandra 调用,所以请查看 DB 命中原因 - 这需要更少的时间
  2. 我发现的违规部分是为每个分区创建 SNS/SQS 连接

它占用了整个工作时间的 60% 以上

我正在 foreachPartition 中创建 SNS/SQS 连接以减少连接。我们有更好的方法吗

我无法在驱动程序上创建连接对象,因为它们不可序列化

我没有使用 executor 9,executore core 15,driver memory 2g,executor memory 5g

我正在使用 16 核 64 gig 内存 集群大小 1 主 9 从 全部相同的配置 EMR 部署火花 1.6

【问题讨论】:

  • 您确定 create an aws SNS and SQS client connection 占用了 60% 的工作时间还是 publish the record to SNS/SQS 这个?这两者之间存在细微差别。对于第一种情况,您需要最小化连接创建的数量,而对于第二种情况,您需要分发数据(并创建更多连接实例)。有趣!!!!
  • 如果是第二种情况,我会发布一个带有解决方案的答案。

标签: apache-spark connection spark-streaming amazon-sns


【解决方案1】:

听起来您希望每个节点只设置一个 SNS/SQS 连接,然后使用它来处理每个节点上的所有数据。

我认为 foreachPartition 在这里是正确的想法,但您可能希望事先合并您的 RDD。这将折叠同一节点上的分区而不进行混洗,并允许您避免启动额外的 SNS/SQS 连接。

请看这里: https://spark.apache.org/docs/latest/api/scala/index.html#org.apache.spark.rdd.RDD@coalesce(numPartitions:Int,shuffle:Boolean,partitionCoalescer:Option[org.apache.spark.rdd.PartitionCoalescer])(implicitord:Ordering[T]):org.apache.spark.rdd.RDD[T]

【讨论】:

  • 是的,coalesce 正是我的解决方案。我想在这里补充一点。我有许多小文件,如 23kb 、 45 kb 等,并且通过 coalesce 调用它缩小到正确的分区,现在我能够在 20 分钟内处理接近 25gb 的文件。在这里进一步改进
  • 谢谢布拉德利..还有一件事..这是说我需要 1TB 数据来处理我应该创建多少个带合并的分区?
  • 所以我会使用一些足够大的分区,以便每个分区都适合内存,或者我拥有的核心数量。以较大者为准。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2020-07-25
  • 1970-01-01
  • 1970-01-01
  • 2015-08-09
  • 1970-01-01
  • 2015-05-12
  • 2019-10-11
相关资源
最近更新 更多