Google Colab环境下PySpark连接Cassandra出现数据源识别错误如何解决?
Google Colab PySpark连接Cassandra解决方案
问题根因
代码出现org.apache.spark.sql.cassandra相关报错的核心原因如下:
- 连接器Maven坐标拼写错误:原代码中
com.datastax.spark:spark-cassandra-connector2.12:3.1.0缺少版本分隔下划线,正确应为com.datastax.spark:spark-cassandra-connector_2.12:3.1.0 - 配置项无效:
os.environ['SPARK_SUBMIT']不是Spark可识别的环境变量,配置不生效 - 依赖拉取方式错误:使用
--jars参数指定Maven坐标无法自动下载依赖,需要改用--packages参数让Spark自动拉取连接器及配套依赖 - 冗余重复配置
SPARK_HOME,且安装了findspark但未调用初始化方法,会导致Spark环境加载异常 - 已废弃的
SQLContext写法兼容性差,容易出现类加载错误
修正后完整可运行代码
# 安装基础依赖 !apt-get install openjdk-8-jdk-headless -qq > /dev/null !wget https://downloads.apache.org/spark/spark-3.1.2/spark-3.1.2-bin-hadoop3.2.tgz !tar -xvzf spark-3.1.2-bin-hadoop3.2.tgz !pip install findspark pyspark==3.1.2 # 配置环境变量 import os os.environ["JAVA_HOME"] = "/usr/lib/jvm/java-8-openjdk-amd64" os.environ["SPARK_HOME"] = "/content/spark-3.1.2-bin-hadoop3.2" # 配置自动拉取Cassandra连接器依赖 os.environ['PYSPARK_SUBMIT_ARGS'] = '--packages com.datastax.spark:spark-cassandra-connector_2.12:3.1.0 pyspark-shell' # 初始化Spark环境 import findspark findspark.init() from pyspark.sql import SparkSession # 初始化Spark会话 spark = SparkSession.builder \ .appName("Spark Cassandra") \ .config("spark.cassandra.connection.host", "你的Cassandra主机地址") \ .config("spark.cassandra.auth.username", "你的用户名") \ .config("spark.cassandra.auth.password", "你的密码") \ .getOrCreate() # 读取Cassandra表数据 dataFrame = spark.read.format("org.apache.spark.sql.cassandra") \ .options(table="你的表名", keyspace="你的键空间名") \ .load() dataFrame.printSchema()
注意事项
- 请替换代码中所有占位的配置项为你自己的Cassandra集群实际信息
- 若你的Cassandra使用非默认9042端口,需要额外添加配置项
.config("spark.cassandra.connection.port", "你的端口号") - 首次运行时Spark会自动下载连接器依赖,耗时1-3分钟属于正常情况,请勿中断执行
- 若依赖拉取失败,可手动下载对应版本的连接器jar包上传到Colab,再将
--packages改为--jars /content/对应jar包文件名.jar指定本地jar路径即可
内容的提问来源于stack exchange,提问作者M.Izzath
相关产品推荐
相关产品推荐

