Spark应用连接Cassandra失败:IOException问题求助
Kafka-Spark Streaming写入Cassandra连接失败问题排查
使用组件版本
- spark-2.1.1-bin-hadoop2.7
- kafka_2.11-0.9.0.0
- apache-cassandra-3.9
编写的代码
val sparkConf = new SparkConf().setAppName("KafkaSparkStreaming").set("spark.cassandra.connection.host", "127.0.0.1") val ssc = new StreamingContext(sparkConf, Seconds(5)) val topicpMap = "mytopic".split(",").map((_, 1.toInt)).toMap val lines = KafkaUtils.createStream(ssc, "localhost:2181", "sparkgroup", topicpMap).map(_._2) lines.map(line => { val arr = line.split(","); (arr(0),arr(1),arr(2),arr(3),arr(4)) }).saveToCassandra("sparkdata", "cust_data", SomeColumns("fname", "lname","url","product","cnt"))
运行报错信息
java.io.IOException: Failed to open native connection to Cassandra at {127.0.0.1}:9042 at com.datastax.spark.connector.cql.CassandraConnector$.com$datastax$spark$connector$cql$CassandraConnector$$createSession(CassandraConnector.scala:168) at com.datastax.spark.connector.cql.CassandraConnector$$anonfun$8.apply(CassandraConnector.scala:154) at com.datastax.spark.connector.cql.CassandraConnector$$anonfun$8.apply(CassandraConnector.scala:154) at com.datastax.spark.connector.cql.RefCountedCache.createNewValueAndKeys(RefCountedCache.scala:32) at com.datastax.spark.connector.cql.RefCountedCache.syncAcquire(RefCountedCache.scala:69) at com.datastax.spark.connector.cql.RefCountedCache.acquire(RefCountedCache.scala:57) at com.datastax.spark.connector.cql.CassandraConnector.openSession(CassandraConnector.scala:79) at com.datastax.spark.connector.cql.CassandraConnector.withSessionDo(CassandraConnector.scala:111) at com.datastax.spark.connector.cql.CassandraConnector.withClusterDo(CassandraConnector.scala:122) at com.datastax.spark.connector.cql.Schema$.fromCassandra(Schema.scala:330) at com.datastax.spark.connector.cql.Schema$.tableFromCassandra(Schema.scala:350) at com.datastax.spark.connector.writer.TableWriter$.apply(TableWriter.scala:336) at com.datastax.spark.connector.streaming.DStreamFunctions.saveToCassandra(DStreamFunctions.scala:53) ... 58 elided Caused by: com.datastax.driver.core.exceptions.NoHostAvailableException: All host(s) tried for query failed (tried: /127.0.0.1:9042 (com.datastax.driver.core.exceptions.TransportException: [/127.0.0.1:9042] Cannot connect)) at com.datastax.driver.core.ControlConnection.reconnectInternal(ControlConnection.java:233) at com.datastax.driver.core.ControlConnection.connect(ControlConnection.java:79) at com.datastax.driver.core.Cluster$Manager.init(Cluster.java:1483) at com.datastax.driver.core.Cluster.getMetadata(Cluster.java:399) at com.datastax.spark.connector.cql.CassandraConnector$.com$datastax$spark$connector$cql$CassandraConnector$$createSession(CassandraConnector.scala:161) ... 70 more
问题原因及解决方法
核心问题
错误栈明确显示Spark无法连接到Cassandra的9042端口(CQL原生协议端口),以下是具体排查方向:
Cassandra服务未启动
- 执行
nodetool status检查节点状态,确认节点处于UP状态; - 若未启动,执行
cassandra -f(前台启动便于查看日志)或对应后台启动命令启动服务。
- 执行
Cassandra网络配置错误
- 打开
cassandra.yaml配置文件,检查以下项:listen_address:设置为Spark可访问的地址,本地测试可设为127.0.0.1或localhost;rpc_address:CQL连接的服务地址,需与Spark配置的spark.cassandra.connection.host一致;native_transport_port:确认端口为9042,若修改过需同步调整Spark配置中的spark.cassandra.connection.port。
- 打开
端口占用或防火墙拦截
- 检查9042端口是否被占用:Linux执行
netstat -an | grep 9042,Windows执行netstat -ano | findstr 9042,若被占用则停止对应进程或修改Cassandra端口; - 测试环境可临时关闭防火墙,生产环境需在防火墙上开放9042端口。
- 检查9042端口是否被占用:Linux执行
Spark-Cassandra Connector版本不兼容
- Spark 2.1.1需搭配2.0.x系列的Connector(如
com.datastax.spark:spark-cassandra-connector_2.11:2.0.10),确认项目依赖中引入了对应版本的Connector,避免因版本不匹配导致连接失败。
- Spark 2.1.1需搭配2.0.x系列的Connector(如
内容的提问来源于stack exchange,提问作者UTHAYAKUMAR M
相关产品推荐
相关产品推荐

