Here is the code I am running on the EMR
import pyspark
from pyspark.sql import SQLContext
from pyspark.sql import SparkSession
from pyspark import SQLContext, SparkContext, SparkConf
spark = SparkSession.builder.getOrCreate()
sql_context = SQLContext(spark)
sc = spark.sparkContext
url = "jdbc:redshift://redshift-cluster-endpoint.amazonaws.com::5439/db-name?user=<USER>&password=<PWD>"
df = sql_context.read \
.format("com.databricks.spark.redshift") \
.option("url", url) \
.option("dbtable", "schema_name.table_name") \
.option("tempdir", "s3://temp-bucket/temp_data_pyspark/") \
.option("forward_spark_s3_credentials", "true") \
.load()
print(df.head(10))
This is my spark-submit
spark-submit \
--jars s3://aws-emr-resources-bucket-us-east-1/RedshiftJDBC41-1.2.12.1017.jar,s3://aws-emr-resources-bucket-us-east-1/minimal-json-0.9.4.jar,s3://aws-emr-resources-bucket-us-east-1/spark-avro_2.11-3.0.0.jar,s3://aws-emr-resources-bucket-us-east-1/spark-redshift_2.10-2.0.0.jar \
--packages com.databricks:spark-redshift_2.11:2.0.1 \
python-script.py
The error I am facing is :
Traceback (most recent call last):
File "/home/hadoop/python-script.py", line 18, in <module>
.option("forward_spark_s3_credentials", "true") \
File "/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/readwriter.py", line 210, in load
File "/usr/lib/spark/python/lib/py4j-0.10.9-src.zip/py4j/java_gateway.py", line 1305, in __call__
File "/usr/lib/spark/python/lib/pyspark.zip/pyspark/sql/utils.py", line 111, in deco
File "/usr/lib/spark/python/lib/py4j-0.10.9-src.zip/py4j/protocol.py", line 328, in get_return_value
py4j.protocol.Py4JJavaError: An error occurred while calling o67.load.
: java.lang.NoClassDefFoundError: scala/Product$class
at com.databricks.spark.redshift.Parameters$MergedParameters.<init>(Parameters.scala:78)
at com.databricks.spark.redshift.Parameters$.mergeParameters(Parameters.scala:72)
at com.databricks.spark.redshift.DefaultSource.createRelation(DefaultSource.scala:48)
at org.apache.spark.sql.execution.datasources.DataSource.resolveRelation(DataSource.scala:355)
at org.apache.spark.sql.DataFrameReader.loadV1Source(DataFrameReader.scala:325)
at org.apache.spark.sql.DataFrameReader.$anonfun$load$3(DataFrameReader.scala:307)
at scala.Option.getOrElse(Option.scala:189)
at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:307)
at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:225)
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:498)
at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357)
at py4j.Gateway.invoke(Gateway.java:282)
at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
at py4j.commands.CallCommand.execute(CallCommand.java:79)
at py4j.GatewayConnection.run(GatewayConnection.java:238)
at java.lang.Thread.run(Thread.java:750)
Caused by: java.lang.ClassNotFoundException: scala.Product$class
at java.net.URLClassLoader.findClass(URLClassLoader.java:387)
at java.lang.ClassLoader.loadClass(ClassLoader.java:418)
at java.lang.ClassLoader.loadClass(ClassLoader.java:351)
... 20 more
22/09/06 08:29:38 INFO SparkContext: Invoking stop() from shutdown hook
22/09/06 08:29:38 INFO AbstractConnector: Stopped Spark@2edd2970{HTTP/1.1, (http/1.1)}{0.0.0.0:4040}
22/09/06 08:29:38 INFO SparkUI: Stopped Spark web UI at http://ip-10-50-105-28.ec2.internal:4040
22/09/06 08:29:38 INFO YarnClientSchedulerBackend: Interrupting monitor thread
22/09/06 08:29:38 INFO YarnClientSchedulerBackend: Shutting down all executors
22/09/06 08:29:38 INFO YarnSchedulerBackend$YarnDriverEndpoint: Asking each executor to shut down
22/09/06 08:29:38 INFO YarnClientSchedulerBackend: YARN client scheduler backend Stopped
22/09/06 08:29:38 INFO MapOutputTrackerMasterEndpoint: MapOutputTrackerMasterEndpoint stopped!
22/09/06 08:29:38 INFO MemoryStore: MemoryStore cleared
22/09/06 08:29:38 INFO BlockManager: BlockManager stopped
22/09/06 08:29:38 INFO BlockManagerMaster: BlockManagerMaster stopped
22/09/06 08:29:38 INFO OutputCommitCoordinator$OutputCommitCoordinatorEndpoint: OutputCommitCoordinator stopped!
22/09/06 08:29:38 INFO SparkContext: Successfully stopped SparkContext
22/09/06 08:29:38 INFO ShutdownHookManager: Shutdown hook called
22/09/06 08:29:38 INFO ShutdownHookManager: Deleting directory /mnt/tmp/spark-fefa4b92-1737-4087-ba47-e8a81a3aa2e0
22/09/06 08:29:38 INFO ShutdownHookManager: Deleting directory /mnt/tmp/spark-a58798f3-021f-4a8a-88f6-5bd54d57b428/pyspark-442297be-93ef-494d-af78-57b88db7395e
22/09/06 08:29:38 INFO ShutdownHookManager: Deleting directory /mnt/tmp/spark-a58798f3-021f-4a8a-88f6-5bd54d57b428
Upon searching this error, I found here that I must be using different Scala version binaries. Any help in identifying the compatible versions of the above jars is appreciated. Resources/Links would be most beneficial.