Spark 4.1.1 on AKS连接Cosmos DB Cassandra API适配方案咨询
Spark 4.1.1 on AKS 连接Azure Cosmos DB Cassandra API 适配方案
原可用配置(Spark 3.5.3)
- 连接器:
com.datastax.spark:spark-cassandra-connector-assembly_2.12:3.5.1 - 目标数据源:Azure Cosmos DB Cassandra API
- 读写方式:使用Cassandra数据源(
org.apache.spark.sql.cassandra)
升级至Spark 4.1.1后的问题
场景A - 未提供连接器jar包
- 调用
.format("org.apache.spark.sql.cassandra")时触发SparkClassNotFoundException
场景B - 提供旧版连接器jar包
- 使用
spark-cassandra-connector-assembly_2.12:3.5.1后,类缺失问题解决,但读写过程中触发java.io.IOException/ClosedConnectionException,连接失败
适配Spark 4.1.1的正确配置
连接器版本选择
Spark 4.1.1基于Scala 2.13,需使用兼容Spark 4.x且对应Scala版本的连接器。推荐使用:com.datastax.spark:spark-cassandra-connector-assembly_2.13:4.1.0
该版本是Datastax官方适配Spark 4.1.x的稳定版本,完全兼容Azure Cosmos DB Cassandra API。
必要配置变更
- Scala版本匹配:必须选择后缀为
_2.13的连接器构件,Spark 4.1.1不再兼容Scala 2.12版本的连接器。 - Cosmos DB专属连接参数:需在读写配置中添加以下强制参数:
spark.cassandra.connection.host:Cosmos DB Cassandra API的接触点地址spark.cassandra.connection.port:固定为10350(Cosmos DB Cassandra API默认端口)spark.cassandra.connection.ssl.enabled:设为true(Cosmos DB强制SSL连接)spark.cassandra.auth.username:Cosmos DB账户名称spark.cassandra.auth.password:Cosmos DB账户密钥
- AKS提交方式:通过
spark-submit提交时,可直接用--packages参数拉取连接器:spark-submit --packages com.datastax.spark:spark-cassandra-connector-assembly_2.13:4.1.0 \ --class <your-main-class> \ <your-jar-path>
修正后的示例代码
val df = spark.read .format("org.apache.spark.sql.cassandra") .options(Map( "keyspace" -> "<your-keyspace>", "table" -> "<your-table>", "spark.cassandra.connection.host" -> "<cosmos-db-contact-point>", "spark.cassandra.connection.port" -> "10350", "spark.cassandra.connection.ssl.enabled" -> "true", "spark.cassandra.auth.username" -> "<cosmos-db-account-name>", "spark.cassandra.auth.password" -> "<cosmos-db-account-key>" )) .load()
环境信息
- 部署环境:AKS(Kubernetes),通过
spark-submit提交Spark作业 - Spark版本:3.5.3(原可用)→ 4.1.1(当前升级目标)
- 目标数据源:Azure Cosmos DB Cassandra API
内容的提问来源于stack exchange,提问作者akshay kadam
相关产品推荐
相关产品推荐

