Spark Cassandra CassandraSourceRelation directJoinSetting exception error

Viewed 124
    

    // Input Identifiers
    val ids = List("4723847392423894", "4329479647236423", "42348726782684")


    import spark.implicits._
    val settings = Map("table" -> "table_name", "keyspace" -> "keyspace_name")
    val tableDF = spark.read.format("org.apache.spark.sql.cassandra").options(settings).load()
    val idsListDF = ids.asInstanceOf[List[String]].toDF("id").persist()
    idsListDF.join(tableDF, tableDF.col("id") === idsListDF.col("id"), "inner").persist()



Exception

Exception in thread "main" java.lang.NoSuchMethodError: org.apache.spark.sql.cassandra.CassandraSourceRelation.directJoinSetting()Lorg/apache/spark/sql/cassandra/DirectJoinSetting;
    at org.apache.spark.sql.cassandra.execution.CassandraDirectJoinStrategy$.containsSafePlans(CassandraDirectJoinStrategy.scala:333)
    at org.apache.spark.sql.cassandra.execution.CassandraDirectJoinStrategy$.validJoinBranch(CassandraDirectJoinStrategy.scala:283)
    at org.apache.spark.sql.cassandra.execution.CassandraDirectJoinStrategy.rightValid(CassandraDirectJoinStrategy.scala:139)
    at org.apache.spark.sql.cassandra.execution.CassandraDirectJoinStrategy.hasValidDirectJoin(CassandraDirectJoinStrategy.scala:87)
    at org.apache.spark.sql.cassandra.execution.CassandraDirectJoinStrategy.apply(CassandraDirectJoinStrategy.scala:30)
    at org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$1.apply(QueryPlanner.scala:63)
    at org.apache.spark.sql.catalyst.planning.QueryPlanner$$anonfun$1.apply(QueryPlanner.scala:63)
    at scala.collection.Iterator$$anon$12.nextCur(Iterator.scala:435)
    at scala.collection.Iterator$$anon$12.hasNext(Iterator.scala:441)
    at scala.collection.Iterator$$anon$12.hasNext(Iterator.scala:440)



Could you please help me with what's wrong with the code?

I have tried the directJoin(Automatic) automatic, always, always off but still no luck

idsListDF.join(tableDF.directJoin(Automatic), tableDF.col("batch_id") === idsListDF.col("id"), "inner").persist()

FYI - I'm using the Spark Cassandra connector jar - https://github.com/datastax/spark-cassandra-connector

1 Answers

This looks like an environmental issue although I'm unable to pinpoint the cause.

I've engaged the Analytics team here at DataStax and I will post an update once I get a response. Cheers!

P.S. Thanks for posting the Spark + connector versions. I'd recommend updating your original question with those details to make it easier for other contributors to help you.

Related