【问题标题】:Spark broadcasting cassandra ConnectorSpark广播cassandra连接器
【发布时间】:2015-10-08 01:06:21
【问题描述】:

我正在使用 datastax 提供的 spark-cassandra-connector 1.1.0。我注意到有趣的问题,但我不确定为什么会发生这样的事情: 当我广播 cassandra 连接器并尝试在执行程序上使用它时,我收到异常提示我的配置无效无法在 0.0.0 连接到 Cassandra。

示例堆栈跟踪:

java.io.IOException: Failed to open native connection to Cassandra at {0.0.0.0}:9042
        at com.datastax.spark.connector.cql.CassandraConnector$.com$datastax$spark$connector$cql$CassandraConnector$$createSession(CassandraConnector.scala:174)
        at com.datastax.spark.connector.cql.CassandraConnector$$anonfun$2.apply(CassandraConnector.scala:160)
        at com.datastax.spark.connector.cql.CassandraConnector$$anonfun$2.apply(CassandraConnector.scala:160)
        at com.datastax.spark.connector.cql.RefCountedCache.createNewValueAndKeys(RefCountedCache.scala:36)
        at com.datastax.spark.connector.cql.RefCountedCache.acquire(RefCountedCache.scala:61)
        at com.datastax.spark.connector.cql.CassandraConnector.openSession(CassandraConnector.scala:71)
        at com.datastax.spark.connector.cql.CassandraConnector.withSessionDo(CassandraConnector.scala:97)
...

但如果我在不广播的情况下使用它,一切正常。

对我来说也很奇怪,在驱动程序端广播值打印正确的配置,但在执行器端没有。

司机端:

  val dbConf = ssc.sparkContext.getConf
  val connector = CassandraConnector(dbConf)
  println(connector.hosts) //Set(10.20.1.5) 
  val broadcastedConnector = ssc.sparkContext.broadcast(connector)
  println(broadcastedConnector.value.hosts) //Set(10.20.1.5) 

执行者端:

mapPartition{
...
 println(broadcastedConnector.hosts) // Set(0.0.0.)
...
}

有人可以解释为什么它以这种方式工作,以及如何以一种可以在执行者方面使用的方式广播 Cassandra 连接器。

更新同样的问题出现在 1.2.3 版本的连接器中。

【问题讨论】:

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


    【解决方案1】:

    没有理由广播 Cassandra 连接器。在并行闭包中使用它只会序列化配置并在执行器上创建新连接,或者使用现有的执行器连接(如果存在)。

    【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2018-11-12
    • 1970-01-01
    • 1970-01-01
    • 2015-03-12
    • 2018-11-28
    • 2016-02-16
    • 1970-01-01
    • 2017-01-13
    相关资源
    最近更新 更多