Scala+DSE SearchAnalytics读取Cassandra表遇表不存在错误求助
解决Spark SQL无法找到Cassandra表的问题
你遇到的核心问题是:Spark默认不会自动识别Cassandra的keyspace和表,必须通过Cassandra Spark Connector完成两者的映射,才能让Spark SQL正确找到目标表。我给你一步步拆解解决方案:
1. 先确保添加了Cassandra Spark Connector依赖
这是Spark和Cassandra通信的基础,没有它俩根本没法互通。根据你的构建工具添加对应依赖:
- SBT:
libraryDependencies += "com.datastax.spark" %% "spark-cassandra-connector" % "3.2.0" // 版本要匹配你的Spark/DSE版本 - Maven:
<dependency> <groupId>com.datastax.spark</groupId> <artifactId>spark-cassandra-connector_2.12</artifactId> <version>3.2.0</version> </dependency>
注意:版本要和你的Spark、DataStax Enterprise(DSE)版本兼容,比如DSE 6.0要选对应适配的连接器版本。
2. 两种正确查询Cassandra表的方式
方式一:直接加载为DataFrame后查询
用cassandraFormat方法把Cassandra表加载成Spark DataFrame,再执行Solr过滤:
import com.datastax.spark.connector._ val spark = SparkSession.builder() .appName("CassandraSpark") .config("spark.cassandra.connection.host", "127.0.0.1") .config("spark.cassandra.connection.port", "9042") .master("local[2]") .getOrCreate() // 加载Cassandra表为DataFrame val videosDF = spark.read .format("org.apache.spark.sql.cassandra") .options(Map( "keyspace" -> "killr_video", "table" -> "videos" )) .load() // 用Solr查询过滤数据 val filteredDF = videosDF.filter("solr_query = '{\"q\":\"video_id:1\"}'") filteredDF.show()
方式二:创建临时视图后用Spark SQL查询
如果你更习惯SQL语法,可以先把Cassandra表注册成Spark临时视图,之后就能像操作普通SQL表一样查询:
import com.datastax.spark.connector._ val spark = SparkSession.builder() .appName("CassandraSpark") .config("spark.cassandra.connection.host", "127.0.0.1") .config("spark.cassandra.connection.port", "9042") .master("local[2]") .getOrCreate() // 注册Cassandra表为临时视图 spark.read .format("org.apache.spark.sql.cassandra") .options(Map( "keyspace" -> "killr_video", "table" -> "videos" )) .createOrReplaceTempView("videos") // 现在可以用Spark SQL执行Solr查询了 val ss = spark.sql("select * from videos where solr_query = '{\"q\":\"video_id:1\"}'") ss.show()
为什么你的原代码会报错?
你直接用spark.sql("select * from killr_video.videos ...")时,Spark会在自身的元数据中查找这个表,但它完全不知道killr_video是Cassandra的keyspace——必须通过连接器把Cassandra表加载到Spark上下文,或者注册成临时视图,Spark才能识别到它的存在。
另外,如果你用的是DSE Search Analytics,只要你的SparkSession已经指定了正确的连接地址,且节点网络连通、状态正常,上面的方法就能正常生效。
内容的提问来源于stack exchange,提问作者Chinmay R
相关产品推荐
相关产品推荐

