Come utilizzare i pacchetti Spark in AWS Glue?
Vorrei utilizzare il connettore spark-cassandra di Datastax in AWS Glue. Se eseguo pyspark localmente, il mio comando sarebbe simile a
path/to/spark-3.0.1-bin-hadoop2.7/bin/spark-submit \
--conf spark.cassandra.connection.host=XXX \
--conf spark.cassandra.auth.username=XXX \
--conf spark.cassandra.auth.password=XXX \
--packages com.datastax.spark:spark-cassandra-connector_2.12:2.5.1 \
~/my_script.py
Come eseguire questo script in Glue?
Cose che ho provato
Come importare pacchetti Spark in AWS Glue? Sembra simile alla mia domanda. La risposta accettata parla dell'aggiunta di un modulo Python zippato come parametro. Ma
spark-cassandra-connectornon è un modulo Python.(secondo il commento di @ alex) metti l'assieme SCC nel lavoro di colla
Jar lib path
Errore:
File "/tmp/delta_on_s3_spark.py", line 75, in _write_df_to_cassandra
df.write.format(format_).mode('append').options(table=table, keyspace=keyspace).save()
File "/opt/amazon/spark/python/lib/pyspark.zip/pyspark/sql/readwriter.py", line 732, in save
self._jwrite.save()
File "/opt/amazon/spark/python/lib/py4j-0.10.7-src.zip/py4j/java_gateway.py", line 1257, in __call__
answer, self.gateway_client, self.target_id, self.name)
File "/opt/amazon/spark/python/lib/pyspark.zip/pyspark/sql/utils.py", line 63, in deco
return f(*a, **kw)
File "/opt/amazon/spark/python/lib/py4j-0.10.7-src.zip/py4j/protocol.py", line 328, in get_return_value
format(target_id, ".", name), value)
py4j.protocol.Py4JJavaError: An error occurred while calling o84.save.
: java.lang.NoSuchMethodError: scala.Product.$init$(Lscala/Product;)V
at com.datastax.spark.connector.TableRef.<init>(TableRef.scala:4)
at org.apache.spark.sql.cassandra.DefaultSource$.TableRefAndOptions(DefaultSource.scala:142)
at org.apache.spark.sql.cassandra.DefaultSource.createRelation(DefaultSource.scala:83)
......
- (secondo il commento di @ alex) ha inserito
spark.jars.packages = com.datastax.spark:spark-cassandra-connector_2.12:2.5.1il lavoro di Collajob parameter
Errore:
File "/tmp/delta_on_s3_spark.py", line 75, in _write_df_to_cassandra
df.write.format(format_).mode('append').options(table=table, keyspace=keyspace).save()
File "/opt/amazon/spark/python/lib/pyspark.zip/pyspark/sql/readwriter.py", line 732, in save
self._jwrite.save()
File "/opt/amazon/spark/python/lib/py4j-0.10.7-src.zip/py4j/java_gateway.py", line 1257, in __call__
answer, self.gateway_client, self.target_id, self.name)
File "/opt/amazon/spark/python/lib/pyspark.zip/pyspark/sql/utils.py", line 63, in deco
return f(*a, **kw)
File "/opt/amazon/spark/python/lib/py4j-0.10.7-src.zip/py4j/protocol.py", line 328, in get_return_value
format(target_id, ".", name), value)
py4j.protocol.Py4JJavaError: An error occurred while calling o83.save.
: java.lang.ClassNotFoundException: Failed to find data source: org.apache.spark.sql.cassandra. Please find packages at http://spark.apache.org/third-party-projects.html
at org.apache.spark.sql.execution.datasources.DataSource$.lookupDataSource(DataSource.scala:657)
at org.apache.spark.sql.DataFrameWriter.save(DataFrameWriter.scala:245)
at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
......
Risposte
Il modo consigliato è usare --packageso --conf spark.jars.packagescon le coordinate Maven, quindi Spark estrarrà correttamente tutte le dipendenze necessarie utilizzate da Spark Cassandra Connector (driver Java, ecc.) - Se usi --jarssolo il jar SCC, il tuo lavoro fallirà.
A partire da SCC 2.5.1, c'è anche un nuovo artefatto: spark-cassandra-connector-assembly che include tutte le dipendenze necessarie. Con esso puoi evitare problemi con dipendenze in conflitto, inoltre puoi usarlo con --jarso con il percorso della libreria Jar di Glue job.
PS Con Spark 3.0 si consiglia di utilizzare SCC 3.0.0-beta, a causa di cambiamenti significativi negli interni di Spark SQL.