【问题标题】:Spark Structured Streaming from kafka to save data in Cassandra in Distributed fashion来自 kafka 的 Spark Structured Streaming 以分布式方式在 Cassandra 中保存数据
【发布时间】:2019-06-28 04:26:43
【问题描述】:

我正在尝试创建从 Kafka 到 Spark 的结构化流,这是一个 json 字符串。现在想要将 json 解析为特定列,然后以最佳速度将数据帧保存到 cassandra 表中。使用 Spark 2.4 和 cassandra 2.11 (Apache) 而不是 DSE。

我尝试创建一个 Direct Stream,它提供案例类的 DStream,我在 DStream 上使用 foreachRDD 将其保存到 Cassandra 中,但这每 6-7 天后就会挂起。因此尝试流式传输直接提供数据帧并可以保存到 Cassandra。

val conf = new SparkConf()
          .setMaster("local[3]")
      .setAppName("Fleet Live Data")
      .set("spark.cassandra.connection.host", "ip")
      .set("spark.cassandra.connection.keep_alive_ms", "20000")
      .set("spark.cassandra.auth.username", "user")
      .set("spark.cassandra.auth.password", "pass")
      .set("spark.streaming.stopGracefullyOnShutdown", "true")
      .set("spark.executor.memory", "2g")
      .set("spark.driver.memory", "2g")
      .set("spark.submit.deployMode", "cluster")
      .set("spark.executor.instances", "4")
      .set("spark.executor.cores", "2")
      .set("spark.cores.max", "9")
      .set("spark.driver.cores", "9")
      .set("spark.speculation", "true")
      .set("spark.locality.wait", "2s")

val spark = SparkSession
  .builder
  .appName("Fleet Live Data")
  .config(conf)
  .getOrCreate()
println("Spark Session Config Done")

val sc = SparkContext.getOrCreate(conf)
sc.setLogLevel("ERROR")
val ssc = new StreamingContext(sc, Seconds(10))
val sqlContext = new SQLContext(sc)
 val topics = Map("livefleet" -> 1)
import spark.implicits._
implicit val formats = DefaultFormats

 val df = spark
  .readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "brokerIP:port")
  .option("subscribe", "livefleet")
  .load()

val collection = df.selectExpr("CAST(value AS STRING)").map(f => parse(f.toString()).extract[liveevent])

val query = collection.writeStream
  .option("checkpointLocation", "/tmp/check_point/")
  .format("kafka")
  .format("org.apache.spark.sql.cassandra")
  .option("keyspace", "trackfleet_db")
  .option("table", "locationinfotemp1")
  .outputMode(OutputMode.Update)
  .start()
  query.awaitTermination()

预期是将数据帧保存到 cassandra。但是收到此错误:-

线程 "main" org.apache.spark.sql.AnalysisException 中的异常:必须使用 writeStream.start() 执行带有流式源的查询

【问题讨论】:

  • 你看过 Kafka Connect 吗?这是 Apache Kafka 的一部分,是一种将数据从 Kafka 主题流式传输到目标数据存储(例如 Cassandra)的好方法。
  • 提示:.format("kafka").format("org.apache.spark.sql.cassandra") 不正确
  • 你有没有在代码末尾调用writeStream.start()
  • @cricket_007 - 我知道它不正确,但我实际上正在寻找解决方案,如果我删除 .format("org.apache.spark.sql.cassandra") 这个应该是正确的部分,然后它可以工作,但在这种情况下,它会开始在控制台上显示并且不会保存到 cassandra。
  • @Pinnacle 是的,Kafka Connect 可以分布式运行。如果您想使用 Spark,那很好,我只是在检查您是否知道可能更适合的替代工具。

标签: scala apache-kafka spark-streaming spark-structured-streaming spark-cassandra-connector


【解决方案1】:

根据错误信息,我会说Cassandra不是Streaming Sink,我相信你需要使用.write

collection.write
    .format("org.apache.spark.sql.cassandra")
    .options(...)
    .save() 

import org.apache.spark.sql.cassandra._

// ...
collection.cassandraFormat(table, keyspace).save()

文档:https://github.com/datastax/spark-cassandra-connector/blob/master/doc/14_data_frames.md#example-using-helper-commands-to-write-datasets


但这可能仅适用于数据帧,对于流式源,请参阅 this example,它使用 .saveToCassandra

import com.datastax.spark.connector.streaming._

// ...
val wc = stream.flatMap(_.split("\\s+"))
    .map(x => (x, 1))
    .reduceByKey(_ + _)
    .saveToCassandra("streaming_test", "words", SomeColumns("word", "count")) 

ssc.start()

如果这不起作用,您确实需要一个 ForEachWriter

collection.writeStream
  .foreach(new ForeachWriter[Row] {

  override def process(row: Row): Unit = {
    println(s"Processing ${row}")
  }

  override def close(errorOrNull: Throwable): Unit = {}

  override def open(partitionId: Long, version: Long): Boolean = {
    true
  }
})
.start()

另外值得一提的是,Datastax 发布了一个 Kafka 连接器,并且 Kafka Connect 包含在您的 Kafka 安装(假设为 0.10.2)或更高版本中。你可以找到它的announcement here

【讨论】:

  • 对于写入策略,我收到此错误 - 线程“主”org.apache.spark.sql.AnalysisException 中的异常:无法在流数据集/数据帧上调用“写入”;在 org.apache.spark.sql.catalyst.analysis.package$AnalysisErrorAt.failAnalysis(package.scala:42)
  • 对于 cassandraformat - 我得到“值 cassandraFormat 不是 org.apache.spark.sql.Dataset[class_name] 的成员”
  • 我想问一下使用 ForEachWriter 是否有额外的开销?
  • .saveToCassandra - 这个策略效果很好,但我实际上只是将 RDD 保存到 DB 中,这再次比 Dataframe 或 Dataset 慢得多
  • 每个.format 内部都使用ForEachWriter。使用.format() 有开销,因为发生了内部序列化。您需要在 Spark 的最基础级别使用 RDD,因为这是在 DataFrame 中指定 Row 的原因。关于第二条评论,请查看文档,并检查导入语句
【解决方案2】:

如果您使用的是 Spark 2.4.0,请尝试使用 foreachbatch 编写器。它在流式查询中使用基于批处理的编写器。

    val query= test.writeStream
       .foreachBatch((batchDF, batchId) =>
        batchDF.write
               .format("org.apache.spark.sql.cassandra")
               .mode(saveMode)
               .options(Map("keyspace" -> keySpace, "table" -> tableName))
               .save())
      .trigger(Trigger.ProcessingTime(3000))
      .option("checkpointLocation", /checkpointing")
      .start
   query.awaitTermination()

【讨论】:

  • 这显示错误无效 saveMode - 允许 append , overwrite 等。接下来我应该用什么替换“test.writestream”,在我的情况下应该是“collection.writestream”或“df.writestream” ” 。如果我采用“collection.writestream”,它会给出“Task not serializable”错误,对于“df.writestream”,它会给出“Columns not found in table db_name.table_name:key、value、topic、partition、offset、timestamp、timestampType”我只想按照问题中描述的模式保存价值。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-05-22
  • 2020-07-25
  • 2021-04-26
  • 2020-12-24
相关资源
最近更新 更多