【问题标题】:can not saveToCassandra in Spark无法在 Spark 中保存 ToCassandra
【发布时间】:2017-10-14 13:09:42
【问题描述】:

我想从 hdfs 获取文件并保存到 cassandra

import org.apache.spark.{SparkConf, SparkContext}
import com.datastax.spark.connector._
val conf = new SparkConf().setMaster("local[2]).setAppName("test")
.set("spark.cassandra.connection.host", "192.168.0.1")
val sc = new SparkContext(conf)
val files = sc.textFiles("hdfs://192.168.0.1:9000/test/", 1)
files.map(_.split("\n")).saveToCassandra("ks", "tb", SomeColumns("id", "time", "text"))
sc.stop()

但由于异常,我无法将其写入 cassandra

我得到的文件,因为 files.foreach(x => println(x)) 有效

【问题讨论】:

  • 错误是什么?
  • 线程“main”java.lang.IllegalArgumentException 中的异常:要求失败:在 scala.Array[String] 中找不到列:[id, time, text]

标签: scala apache-spark hdfs spark-streaming spark-cassandra-connector


【解决方案1】:

据我所知,您可以进行以下更改

var wc1 = files.map(_.split("\\|")).map(r=>Row(r(0),r(1),r(2)))
implicit val rowWriter = SqlRowWriter.Factory
wc1.saveToCassandra("ks", "tb", SomeColumns("id", "time", "text"))

希望它有效!!!...

【讨论】:

  • java.lang.ArrayIndexOutOfBoundsException: 1
  • 您需要转换 files.map(_.split("\n")) github.com/datastax/spark-cassandra-connector/blob/master/doc/…
【解决方案2】:

class randomString4File { val x = util.Random def create(size: Int) = { var str = "" while(str.length() < size) { if (str.length() != 0) {

   str += "\n" + x.nextInt().abs.toString() + " " + java.time.LocalDateTime.now() + " " + x.alphanumeric.take(20).mkString 
  }
  else str += x.nextInt().abs.toString() + " " + java.time.LocalDateTime.now() + " " + x.alphanumeric.take(20).mkString + "\n"
}
str

}

} 我这样写文件

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-01-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多