Spark Structured Stream对接Cassandra查询偶发UnavailableException问题咨询
问题分析与解决方案
一、UnavailableException 错误原因及修复
这个错误表明查询时,Cassandra无法找到满足一致性级别要求的存活副本。结合你的1主3从架构,主要原因及解决办法如下:
1. Spark提交命令的配置错误
你的spark-submit命令中存在两处关键错误:
- 配置项拼写错误:
com.datastax.spark:spark.cassandra.connectiohost应为spark.cassandra.connection.host(多了一个io) - 多个
--conf参数不应通过逗号合并,需单独指定
修正后的spark-submit命令:
spark-submit --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.1.2,com.datastax.spark:spark-cassandra-connector_2.12:3.2.0,com.github.jnr:jnr-posix:3.1.15 \ --conf spark.dynamicAllocation.enabled=false \ --conf spark.cassandra.connection.host=172.16.10.6 \ --conf spark.sql.extensions=com.datastax.spark.connector.CassandraSparkExtensions
2. Cassandra副本策略与一致性级别不匹配
- 检查keyspace的复制因子:如果复制因子过低(比如为1),主节点宕机时会导致无存活副本。执行CQL命令查看:
DESCRIBE KEYSPACE your_keyspace; - 若复制因子不足,调整为合适的值(比如1主3从架构可设为3):
ALTER KEYSPACE your_keyspace WITH REPLICATION = {'class': 'SimpleStrategy', 'replication_factor': 3}; - 查询时的一致性级别:确保查询使用的一致性级别与复制因子匹配。例如使用
LOCAL_ONE时,只要有一个副本存活即可返回结果,但如果对应token范围的所有副本都不可达,就会触发错误。可根据业务需求调整为LOCAL_QUORUM(需要多数副本存活)或降低一致性级别。
3. 节点状态与网络问题
- 用
nodetool status检查Cassandra所有节点是否处于UP状态,若有节点DOWN或UN,排查节点宕机、防火墙(开放9042端口)或跨节点通信(7000端口)问题。 - 不要仅配置单个Cassandra节点地址到
spark.cassandra.connection.host,应加入所有节点IP(如172.16.10.6,172.16.10.7,172.16.10.8,172.16.10.9),让Spark自动切换可用节点。
二、Spark与Cassandra连接自动断开的原因及修复
1. 连接池参数配置不合理
默认连接池参数可能导致空闲连接被Cassandra端关闭,或超时未重连。在SparkSession配置中添加以下参数:
spark = SparkSession.builder \ .master(_connSession)\ .appName("Spark_Streaming+kafka+cassandra") \ .config('spark.cassandra.connection.host', _connCassandraHost) \ .config('spark.cassandra.connection.port', _connCassandraPort) \ # 连接超时时间(毫秒) .config('spark.cassandra.connection.timeout_ms', 10000) \ # 读取请求超时时间 .config('spark.cassandra.read.timeout_ms', 20000) \ # 连接保活时长,避免空闲连接被关闭 .config('spark.cassandra.connection.keep_alive_ms', 300000) \ # 重连延迟时间 .config('spark.cassandra.connection.reconnect_delay_ms', 1000) \ .getOrCreate()
2. 网络环境不稳定
- 确保Spark集群与Cassandra集群之间网络无频繁丢包或中断,开启TCP keepalive维持长连接。
- 检查Cassandra节点的
cassandra.yaml配置,开启TCP保活:tcp_keepalive: true
3. Cassandra端连接超时设置
调整Cassandra的连接超时参数(cassandra.yaml):
connection_timeout_in_ms: 5000 read_request_timeout_in_ms: 10000
额外建议
- 将Structured Streaming的
checkpointLocation改为分布式存储路径(如HDFS),避免本地路径在executor重启时导致的状态丢失问题。 - 监控Cassandra节点的负载和状态,及时发现节点异常。
内容的提问来源于stack exchange,提问作者hi-inbeom
相关产品推荐
相关产品推荐

