You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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原生协议端口),以下是具体排查方向:

  1. Cassandra服务未启动

    • 执行nodetool status检查节点状态,确认节点处于UP状态;
    • 若未启动,执行cassandra -f(前台启动便于查看日志)或对应后台启动命令启动服务。
  2. 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。
  3. 端口占用或防火墙拦截

    • 检查9042端口是否被占用:Linux执行netstat -an | grep 9042,Windows执行netstat -ano | findstr 9042,若被占用则停止对应进程或修改Cassandra端口;
    • 测试环境可临时关闭防火墙,生产环境需在防火墙上开放9042端口。
  4. Spark-Cassandra Connector版本不兼容

    • Spark 2.1.1需搭配2.0.x系列的Connector(如com.datastax.spark:spark-cassandra-connector_2.11:2.0.10),确认项目依赖中引入了对应版本的Connector,避免因版本不匹配导致连接失败。

内容的提问来源于stack exchange,提问作者UTHAYAKUMAR M

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.19 00:28:11