【问题标题】:Spark Pooling Options火花池选项
【发布时间】:2019-01-02 23:03:18
【问题描述】:

我无法理解如何使用池选项的设置,也无法判断它们是否从以下来源工作:https://docs.datastax.com/en/developer/java-driver/3.4/manual/pooling/

SparkSession val 是否会考虑集群中的池选项?

我的 Scala 代码:

package com.zeropoints.processing

import org.apache.spark.SparkContext
import org.apache.spark.SparkConf
import org.apache.spark.sql.SparkSession
import com.datastax.spark.connector._
import com.datastax.driver.core.Cluster
import com.datastax.driver.core.PoolingOptions
import com.datastax.driver.core.HostDistance

//This object provides the main entry point into spark processing
object main {

  var appName = "Processing" 

  lazy val sparkconf:SparkConf = new SparkConf(true).setAppName(appName)

  lazy val poolingOptions:PoolingOptions = new PoolingOptions()

  lazy val cluster:Cluster = Cluster.builder().withPoolingOptions(poolingOptions).build()

  lazy val spark:SparkSession = SparkSession.builder().config(sparkconf).getOrCreate

  lazy val sc:SparkContext = spark.sparkContext

  def main(args: Array[String]) {
    //Set pooling stuff
    poolingOptions.setConnectionsPerHost(HostDistance.LOCAL, 6, 60)

    //DF and RDDs tasks...
    spark.sql("select * from data.raw").groupBy("key1,key2").agg(sum("views")).
        write.format("org.apache.spark.sql.cassandra").options(Map( "table" -> "summary", "keyspace" -> "data")).
        mode(org.apache.spark.sql.SaveMode.Append).save()

    //..more stuff
  }

}

【问题讨论】:

    标签: apache-spark cassandra spark-cassandra-connector


    【解决方案1】:

    不,Spark 连接器不会考虑您的池配置 - 它的工作方式不同,特别是如果您考虑在分布式环境中执行代码 - 您的 setConnectionsPerHost 仅在驱动程序中执行,并且不会不会影响执行者。

    正确的方法是通过 Spark 配置参数指定必要的设置。文档有一个separate section on connection parametersconnection.connections_per_executor_max 可能是您需要的。您还可以编写自己的类来实现 trait CassandraConnectionFactory 并提供 createCluster 函数的实现。然后可以将这个类名指定为connection.factory配置参数。

    但主要问题是 - 您真的需要调整这些选项吗?你觉得处理速度慢吗? Java 驱动程序文档建议每个主机有 1 个连接,以避免给 Cassandra 带来额外的负载。

    【讨论】:

    • 完全同意,增加每台主机的连接数不太可能提高性能。在测试中,我发现您可以获得一些吞吐量改进,从而增加该值,但这通常只有在您完全停止加载时才会如此,即使这样,改进也非常微不足道。我会考虑更新 datastax 驱动程序文档以使其更清楚。
    • 好的,您的 Spark 代码似乎正在达到默认情况下限制为 256 个进行中请求的“远程 DC”。您可以尝试将 connection.local_dc 参数设置为“本地”Casasndra 数据中心的名称 - 在某些情况下,驱动程序无法猜测您的代码所在的位置。
    • 抱歉回复晚了。我会尝试设置它,看看效果如何。感谢您的帮助
    猜你喜欢
    • 1970-01-01
    • 2015-08-07
    • 1970-01-01
    • 1970-01-01
    • 2011-01-22
    • 2017-06-24
    • 2017-02-22
    • 1970-01-01
    • 2011-09-26
    相关资源
    最近更新 更多