无法从pyspark从cassandra数据库加载信息

问题描述

我有代码

import os
from pyspark import SparkContext,SparkFiles,sqlContext,SparkFiles
from pyspark.sql import sqlContext,SparkSession
from pyspark.sql.functions import col

secure_bundle_file=os.getcwd()+'\\secure-connect-dbtest.zip'
sparkSession =SparkSession.builder.appName('SparkCassandraApp')\
  .config('spark.cassandra.connection.config.cloud.path',secure_bundle_file)\
  .config('spark.cassandra.auth.username','test')\
  .config('spark.cassandra.auth.password','testquart')\
  .config('spark.dse.continuousPagingEnabled',False)\
  .master('local[*]').getorCreate()

data = sparkSession.read.format("org.apache.spark.sql.cassandra")\
  .options(table="tbthesis",keyspace="test").load()
data.count()

我想做的是连接到数据库并检索我的数据。该代码可以很好地连接到数据库,但是一旦到达读取行,它就会说:

Exception has occurred: Py4JJavaError
An error occurred while calling o48.load.
: 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:674)
at org.apache.spark.sql.execution.datasources.DataSource$.lookupDataSourceV2(DataSource.scala:728)
at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:230)
at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:203)
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)

有人可以帮我吗?

此外,我想添加有关此代码的更多详细信息:

我想做的是测试从数据库中读取200万条记录所花费的时间,正常的python-cassandra驱动程序在大约1个小时内(使用SimpleStatement)读取了200万条记录,因此我想知道如何使用那200万条记录使用spark可以持续很多时间。

谢谢

解决方法

您的类路径中没有Spark Cassandra Connector软件包,因此找不到相应的类。

您需要使用spark-submit开始工作(pyspark--packages com.datastax.spark:spark-cassandra-connector_2.11:2.5.1)。

如果您真的只想从python代码中执行此操作,则可以在创建.config("spark.jars.packages","com.datastax.spark:spark-cassandra-connector_2.11:2.5.1")时尝试添加SparkSession,但如果类路径已被实例化,则可能并不总是有效。

P.S。即使在本地模式下,Spark通常也应该胜过SimpleStatement,尽管Spark在分布式模式下确实很出色。您确实不应该使用SimpleStatement来重复执行仅在参数上不同的查询-您应该使用prepared statements。请阅读Developing applications with DataStax drivers指南。 DataStax还赠送了Cassandra. The Definitive Guide这本书的第三版-刚出版时就读-我建议阅读。

,

我的问题解决了。

问题不是java,hadoop或spark,不是连接器的下载过程,但我无法下载任何内容,因为我用于此jar的缓存文件夹上有东西。

spark下载外部jar的文件夹是C:\ Users \ UlysesRico.ivy2 \ jars缓存是C:\ Users \ UlysesRico.ivy2 \ cache

我只是删除了缓存和罐子折叠,然后我做到了:

pyspark-打包com.datastax.spark:spark-cassandra-connector_2.11:2.5.1 而且,我下载了所有的jar并为其缓存信息。

问题终于解决了。