PySpark调用df.show()读Cassandra报system.size_estimates无SELECT权限异常
问题原因
报错和业务表读取逻辑本身无关,核心触发逻辑如下:
- Spark DataFrame 采用懒执行机制:
load()方法仅构建逻辑执行计划,不会实际发起Cassandra数据查询,只有执行show()、count()这类action算子时,才会真正启动数据读取流程。 - Spark Cassandra连接器(简称SCC)在分布式读取前,默认会查询Cassandra的
system.size_estimates系统表,获取各token范围的实际存储数据量,以此合理切分Spark读取分片、设置读取并行度,实现最优的分布式读取性能。当前使用的my_user账号仅拥有业务键空间my_keyspace的操作权限,没有该系统表的SELECT权限,因此触发UnauthorizedException报错。
cqlsh操作无异常的原因很简单:cqlsh是单客户端查询工具,不需要做分布式并行读取的分片规划,执行常规增删改查时不会访问system.size_estimates系统表,因此仅拥有业务键空间权限即可正常操作。
解决方案
根据实际权限情况二选一即可:
- 方案一(生产环境推荐,性能最优)
联系Cassandra集群管理员,使用超级账号执行CQL语句给业务账号授予系统表只读权限,授权语句如下:
授权后无需修改原有Spark代码,连接器可基于真实存储统计信息切分读取分片,并行度设置最合理,读取性能最优。GRANT SELECT ON TABLE system.size_estimates TO my_user; - 方案二(无系统表权限申请通道时使用)
如果无法获取系统表查询权限,可在读取配置中新增参数,关闭基于system.size_estimates的分片估算逻辑,连接器会直接基于全量token范围切分分片,不再访问该系统表。修改后的读取代码如下:
注意:该方案下分片切分没有参考真实存储数据量,可能出现部分分片数据倾斜、空分片的问题,读取性能略低于方案一,但可在无系统表权限的场景下正常完成数据读取。df = spark.read \ .format("org.apache.spark.sql.cassandra") \ .option("spark.cassandra.connection.host", "my_host") \ .option("spark.cassandra.connection.port", "9042") \ .option("spark.cassandra.auth.username", "my_user") \ .option("spark.cassandra.auth.password", "my_pass") \ .option("keyspace", "my_keyspace") \ .option("table", "my_table") \ # 关闭size_estimates系统表访问 .option("spark.cassandra.input.size_estimates.enabled", "false") \ # 可根据集群资源情况手动调整单分片大小,单位为MB,默认值64 .option("spark.cassandra.input.split.size_in_mb", "64") \ .load()
内容的提问来源于stack exchange,提问作者user2265417
相关产品推荐
相关产品推荐

