Recently I have upgraded cassandra from v3.0.2 to 3.0.6. The previous version worked fine with Spark (python scripts). Since the upgrade I can't run the scripts that worked before and don't how to fix it.
Below is the code that ran without problems before:
import os
import sys
from os.path import join
from os import getenv
import matplotlib
import matplotlib.pyplot as plt
import math
from datetime import datetime
import numpy as np
import statsmodels.api as sm
APP_NAME = "Traffic Visualization and Forecasting"
# # Path for spark source folder
HOME = '/usr/local'
SPARK_HOME = join(HOME, "spark-1.6.0-bin-hadoop2.6")
os.environ['SPARK_HOME']=SPARK_HOME
PYSPARK_DIR = join(getenv('SPARK_HOME'), 'python')
LIBS = ['py4j-0.9-src.zip']
paths = [join(PYSPARK_DIR, 'lib', libFile) for libFile in LIBS]
paths.append(PYSPARK_DIR)
for path in paths:
if path not in sys.path:
sys.path.insert(1, path)
try:
from pyspark import SparkContext
from pyspark import SparkConf
from pyspark.sql import SQLContext
from pyspark.mllib.stat import Statistics
import pyspark.sql.functions as func
from pyspark.sql.functions import udf, log, col
from pyspark.sql.types import *
from pyspark.sql.window import Window
print ("Successfully imported Spark Modules")
except ImportError as e:
print ("Can not import Spark Modules", e)
sys.exit(1)
conf = SparkConf()
.setAppName(APP_NAME)
.setMaster("spark://Nikos-MacBook-Pro.local:7077")
.set("spark.cassandra.coection.host", "127.0.0.1")
.set("spark.cassandra.auth.useame", "user")
.set("spark.cassandra.auth.password", "test")
sc = SparkContext(conf=conf)
log4jLogger = sc._jvm.org.apache.log4j
logger = log4jLogger.LogManager.getLogger(__name__)
sqlContext = SQLContext(sc)
dfNodes = sqlContext.read.format("org.apache.spark.sql.cassandra")
.options(table="nodes", keyspace="test").load()
dfNodes.describe().show()
Now if I try to run the script it throws the following error:
Py4JJavaError Traceback (most recent call last) in () ----> 1 dfNodes = sqlContext.read.format("org.apache.spark.sql.cassandra").options(table="nodes", keyspace="production").load() 2 dfNodes.describe().show()
/usr/local/spark-1.6.0-bin-hadoop2.6/python/pyspark/sql/readwriter.pyc in load(self, path, format, schema, **options) 137 retu self._df(self._jreader.load(path)) 138 else: --> 139 retu self._df(self._jreader.load()) 140 141 @since(1.4)
/usr/local/spark-1.6.0-bin-hadoop2.6/python/lib/py4j-0.9-src.zip/py4j/java_gateway.py in call(self, *args) 811 answer = self.gateway_client.send_command(command) 812 retu_value = get_retu_value( --> 813 answer, self.gateway_client, self.target_id, self.name) 814 815 for temp_arg in temp_args:
/usr/local/spark-1.6.0-bin-hadoop2.6/python/pyspark/sql/utils.pyc in deco(*a, **kw) 43 def deco(*a, **kw): 44 try: ---> 45 retu f(*a, **kw) 46 except py4j.protocol.Py4JJavaError as e: 47 s = e.java_exception.toString()
/usr/local/spark-1.6.0-bin-hadoop2.6/python/lib/py4j-0.9-src.zip/py4j/protocol.py in get_retu_value(answer, gateway_client, target_id, name) 306 raise Py4JJavaError( 307 "An error occurred while calling {0}{1}{2}.n". --> 308 format(target_id, ".", name), value) 309 else: 310 raise Py4JError(
Py4JJavaError: An error occurred while calling o30.load. : java.lang.ClassNotFoundException: Failed to find data source: org.apache.spark.sql.cassandra. Please find packages at http://spark-packages.org at org.apache.spark.sql.execution.datasources.ResolvedDataSource$.lookupDataSource(ResolvedDataSource.scala:77) at org.apache.spark.sql.execution.datasources.ResolvedDataSource$.apply(ResolvedDataSource.scala:102) at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:119) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:497) at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:231) at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:381) at py4j.Gateway.invoke(Gateway.java:259) at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:133) at py4j.commands.CallCommand.execute(CallCommand.java:79) at py4j.GatewayCoection.run(GatewayCoection.java:209) at java.lang.Thread.run(Thread.java:745) Caused by: java.lang.ClassNotFoundException: org.apache.spark.sql.cassandra.DefaultSource at java.net.URLClassLoader.findClass(URLClassLoader.java:381) at java.lang.ClassLoader.loadClass(ClassLoader.java:424) at java.lang.ClassLoader.loadClass(ClassLoader.java:357) at org.apache.spark.sql.execution.datasources.ResolvedDataSource$$anonfun$4$$anonfun$apply$1.apply(ResolvedDataSource.scala:62) at org.apache.spark.sql.execution.datasources.ResolvedDataSource$$anonfun$4$$anonfun$apply$1.apply(ResolvedDataSource.scala:62) at scala.util.Try$.apply(Try.scala:161) at org.apache.spark.sql.execution.datasources.ResolvedDataSource$$anonfun$4.apply(ResolvedDataSource.scala:62) at org.apache.spark.sql.execution.datasources.ResolvedDataSource$$anonfun$4.apply(ResolvedDataSource.scala:62) at scala.util.Try.orElse(Try.scala:82) at org.apache.spark.sql.execution.datasources.ResolvedDataSource$.lookupDataSource(ResolvedDataSource.scala:62) ... 13 more
In spark-env.sh I have included the following files:
SPARK_CLASSPATH=${SPARK_CLASSPATH}:/usr/lib/apache-cassandra-3.0.6/lib/cassandra-driver-core-3.0.0-shaded.jar
SPARK_CLASSPATH=${SPARK_CLASSPATH}:/usr/lib/apache-cassandra-3.0.6/lib/guava-18.0.jar
# Shared jars
SPARK_CLASSPATH=${SPARK_CLASSPATH}:/usr/lib/spark-plugins/config-1.3.0.jar
SPARK_CLASSPATH=${SPARK_CLASSPATH}:/usr/lib/spark-plugins/fluent-logger-0.2.11.jar
SPARK_CLASSPATH=${SPARK_CLASSPATH}:/usr/lib/spark-plugins/fluent-logger-scala_2.10-0.5.1.jar
SPARK_CLASSPATH=${SPARK_CLASSPATH}:/usr/lib/spark-plugins/json4s-native_2.10-3.3.0.jar
SPARK_CLASSPATH=${SPARK_CLASSPATH}:/usr/lib/spark-plugins/jsr166e-1.1.0.jar
SPARK_CLASSPATH=${SPARK_CLASSPATH}:/usr/lib/spark-plugins/lift-json_2.10-2.6.2.jar
SPARK_CLASSPATH=${SPARK_CLASSPATH}:/usr/lib/spark-plugins/spark-sql_2.10-1.0.0.jar
# Streaming jars
SPARK_CLASSPATH=${SPARK_CLASSPATH}:/usr/lib/spark-plugins/metrics-core-2.2.0.jar
SPARK_CLASSPATH=${SPARK_CLASSPATH}:/usr/lib/spark-plugins/scala-logging-slf4j_2.10-2.1.2.jar
SPARK_CLASSPATH=${SPARK_CLASSPATH}:/usr/lib/spark-plugins/spark-cassandra-coector-1.6.0-M2-s_2.10.jar
With previous version (3.0.2) I have used the following coector:
spark-cassandra-coector_2.10-1.5.0-M3.jar
If anyone knows what is causing the problem and how to fix it I would be very thankful.
