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

如何解决无法从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并为其缓存信息。

问题终于解决了。

版权声明:本文内容由互联网用户自发贡献,该文观点与技术仅代表作者本人。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如发现本站有涉嫌侵权/违法违规的内容, 请发送邮件至 dio@foxmail.com 举报,一经查实,本站将立刻删除。

相关推荐


依赖报错 idea导入项目后依赖报错,解决方案:https://blog.csdn.net/weixin_42420249/article/details/81191861 依赖版本报错:更换其他版本 无法下载依赖可参考:https://blog.csdn.net/weixin_42628809/a
错误1:代码生成器依赖和mybatis依赖冲突 启动项目时报错如下 2021-12-03 13:33:33.927 ERROR 7228 [ main] o.s.b.d.LoggingFailureAnalysisReporter : *************************** APPL
错误1:gradle项目控制台输出为乱码 # 解决方案:https://blog.csdn.net/weixin_43501566/article/details/112482302 # 在gradle-wrapper.properties 添加以下内容 org.gradle.jvmargs=-Df
错误还原:在查询的过程中,传入的workType为0时,该条件不起作用 <select id="xxx"> SELECT di.id, di.name, di.work_type, di.updated... <where> <if test=&qu
报错如下,gcc版本太低 ^ server.c:5346:31: 错误:‘struct redisServer’没有名为‘server_cpulist’的成员 redisSetCpuAffinity(server.server_cpulist); ^ server.c: 在函数‘hasActiveC
解决方案1 1、改项目中.idea/workspace.xml配置文件,增加dynamic.classpath参数 2、搜索PropertiesComponent,添加如下 <property name="dynamic.classpath" value="tru
删除根组件app.vue中的默认代码后报错:Module Error (from ./node_modules/eslint-loader/index.js): 解决方案:关闭ESlint代码检测,在项目根目录创建vue.config.js,在文件中添加 module.exports = { lin
查看spark默认的python版本 [root@master day27]# pyspark /home/software/spark-2.3.4-bin-hadoop2.7/conf/spark-env.sh: line 2: /usr/local/hadoop/bin/hadoop: No s
使用本地python环境可以成功执行 import pandas as pd import matplotlib.pyplot as plt # 设置字体 plt.rcParams['font.sans-serif'] = ['SimHei'] # 能正确显示负号 p
错误1:Request method ‘DELETE‘ not supported 错误还原:controller层有一个接口,访问该接口时报错:Request method ‘DELETE‘ not supported 错误原因:没有接收到前端传入的参数,修改为如下 参考 错误2:cannot r
错误1:启动docker镜像时报错:Error response from daemon: driver failed programming external connectivity on endpoint quirky_allen 解决方法:重启docker -> systemctl r
错误1:private field ‘xxx‘ is never assigned 按Altʾnter快捷键,选择第2项 参考:https://blog.csdn.net/shi_hong_fei_hei/article/details/88814070 错误2:启动时报错,不能找到主启动类 #
报错如下,通过源不能下载,最后警告pip需升级版本 Requirement already satisfied: pip in c:\users\ychen\appdata\local\programs\python\python310\lib\site-packages (22.0.4) Coll
错误1:maven打包报错 错误还原:使用maven打包项目时报错如下 [ERROR] Failed to execute goal org.apache.maven.plugins:maven-resources-plugin:3.2.0:resources (default-resources)
错误1:服务调用时报错 服务消费者模块assess通过openFeign调用服务提供者模块hires 如下为服务提供者模块hires的控制层接口 @RestController @RequestMapping("/hires") public class FeignControl
错误1:运行项目后报如下错误 解决方案 报错2:Failed to execute goal org.apache.maven.plugins:maven-compiler-plugin:3.8.1:compile (default-compile) on project sb 解决方案:在pom.
参考 错误原因 过滤器或拦截器在生效时,redisTemplate还没有注入 解决方案:在注入容器时就生效 @Component //项目运行时就注入Spring容器 public class RedisBean { @Resource private RedisTemplate<String
使用vite构建项目报错 C:\Users\ychen\work>npm init @vitejs/app @vitejs/create-app is deprecated, use npm init vite instead C:\Users\ychen\AppData\Local\npm-