【问题标题】:How to save spark streaming data in cassandra如何在 cassandra 中保存火花流数据
【发布时间】:2018-02-11 11:36:36
【问题描述】:

build.sbt 下面是 build.sbt 文件中包含的内容


val sparkVersion = "1.6.3" 
scalaVersion := "2.10.5" 
resolvers += "Spark Packages Repo" at "https://dl.bintray.com/spark-packages/maven" 
libraryDependencies ++= Seq(
"org.apache.spark" %% "spark-streaming" % sparkVersion,
"org.apache.spark" %% "spark-streaming-kafka" % sparkVersion)
libraryDependencies +="datastax" % "spark-cassandra-connector" % "1.6.3-s_2.10"
libraryDependencies +="org.apache.spark" %% "spark-sql" % "1.1.0"

初始化shell的命令: 下面的命令是我遵循的shell初始化过程

/usr/hdp/2.6.0.3-8/spark/bin/spark-shell --packages datastax:spark-cassandra-connector:1.6.3-s_2.10 --conf spark.cassandra.connection.host=127.0.0.1 –jars spark-streaming-kafka-assembly_2.10-1.6.3.jar

注意: 在这里我专门指定了 jar,因为 SBT 无法获取在后面部分创建 kafkaStream 时使用的 Spark Streaming kafka 所需的库

导入所需的库:

本节包括要在 REPL 会话的各种情况下使用的要导入的库

import org.apache.spark.SparkConf; import org.apache.spark.streaming.StreamingContext; import org.apache.spark.streaming.Seconds; import org.apache.spark.streaming.kafka.KafkaUtils; import com.datastax.spark.connector._ ; import org.apache.spark.sql.cassandra._ ;

设置 Spark Streaming 配置:

这里正在配置火花流上下文所需的配置


val conf = new SparkConf().setMaster("local[*]").setAppName("KafkaReceiver")
conf.set("spark.driver.allowMultipleContexts", "true"); // Required to set this to true because during // shell initialization or starting we a spark context is created with configurations of highlighted
conf.setMaster("local"); // then we are assigning those cofigurations locally

使用上述配置创建 SparkStreamingContext: 使用上面定义的配置,我们以以下方式创建火花流上下文

val ssc = new StreamingContext(conf, Seconds(1)); // Seconds here describe the interval to fetch

使用上面的 Spark Streaming Context aka SSC 创建一个 Kafka 流: 这里 ssc 是上面创建的火花流上下文, “localhost:2181”是 ZKquoram “spark-streaming-consumer-group”是消费者组 Map("test3" -> 5) 是 Map("topic" -> 分区数)

val kafkaStream = KafkaUtils.createStream(ssc, "localhost:2181","spark-streaming-consumer-group", Map("test3" -> 5)).map(_._2)

注意 使用 kafkaStream.print() 打印 kafkaStream 对象时获取的值如下图所示


85052,19,960.00,0,2017-08-29 14:52:41,17,VISHAL_GWY01_HT1,26,VISHAL_GTWY17_PRES_01,1,2,4                                                             
85053,19,167.00,0,2017-08-29 14:52:41,17,VISHAL_GWY01_HT1,25,VISHAL_GTWY1_Temp_01,1,2,4                                                              
85054,19,960.00,0,2017-08-29 14:52:41,17,VISHAL_GWY01_HT1,26,VISHAL_GTWY17_PRES_01,1,2,4                                                             
85055,19,167.00,0,2017-08-29 14:52:54,17,VISHAL_GWY01_HT1,25,VISHAL_GTWY1_Temp_01,1,2,4                                                              
85056,19,960.00,0,2017-08-29 14:52:54,17,VISHAL_GWY01_HT1,26,VISHAL_GTWY17_PRES_01,1,2,4                                                             
85057,19,167.00,0,2017-08-29 14:52:55,17,VISHAL_GWY01_HT1,25,VISHAL_GTWY1_Temp_01,1,2,4                                                              
85058,19,960.00,0,2017-08-29 14:52:55,17,VISHAL_GWY01_HT1,26,VISHAL_GTWY17_PRES_01,1,2,4                                                             

17/09/02 18:25:25 INFO JobScheduler: Finished job streaming job 1504376716000 ms.0 from job set of time 1504376716000 ms                             
17/09/02 18:25:25 INFO JobScheduler: Total delay: 9.661 s for time 1504376716000 ms (execution: 0.021 s)                                             
17/09/02 18:25:25 INFO JobScheduler: Starting job streaming job 1504376717000 ms.0 from job set of time 1504376717000 ms

转换 kafkaStream 并保存在 Cassandra 中:


kafkaStream.foreachRDD( rdd => { 
if (! rdd.isEmpty()) { 
rdd.map( line => { 
val arr = line.split(",");
(arr(0), arr(1), arr(2), arr(3), arr(4), arr(5), arr(6), arr(7), arr(8), arr(9), arr(10), arr(11))
}). saveToCassandra("test", "sensorfeedVals", SomeColumns(
"tableid", "ccid", "paramval", "batVal", "time", "gwid", "gwhName", "snid", "snhName", "snStatus", "sd", "MId")
)
} else {
 println("No records to save")
}
}
)

启动 ssc:

使用 ssc.start 您可以开始流式传输

这里面临的问题是: 1. 只有在我输入exitCtrl+C后才会打印流的内容 2. 每当我使用 ssc.start 时,它会立即在 REPL 中开始流式传输吗?没有给时间输入 ssc.awaitTermination 3. 我在下面的过程中尝试正常保存时的主要问题***

val collection = sc.parallelize(Seq(("key3", 3), ("key4", 4)))
collection.saveToCassandra("test", "kv", SomeColumns("key", "value"))

我能够保存在 Cassandra 中,但每当我尝试使用 转换 kafkaStream 并保存在 Cassandra 中的逻辑尝试保存在 Cassandra 中时: 我无法从字符串中提取每个值并将其保存在Cassandra 表的各个列!

【问题讨论】:

  • “您无法保存这些值”是什么意思,有例外吗?没有显示值吗?您确定有要保存的值吗?
  • @RussS 用我使用 kafkaStream.print() 打印时显示的值编辑了问题。
  • @RussS“我无法保存这些值”意味着使用 kafkaStream 获取的值,我无法遍历这些值并提取 csv 中记录的每个单独值,然后将它们保存到卡桑德拉。
  • “你不能迭代它们”是什么意思。我还是不明白。就像没有写入记录并且当您从 Cassandra 读取记录时没有显示一样?还是抛出异常阻止您进行迭代?
  • @RussS 当我尝试使用上述帖子的“转换 kafkaStream 并保存在 Cassandra”部分中指定的逻辑来迭代该 rdd 时,会引发异常!声明“线程“streaming-job-executor-5”中的异常 java.lang.NoClassDefFoundError:无法初始化类 com.datastax.spark.connector.cql.CassandraC onnector”

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


【解决方案1】:

java.lang.NoClassDefFoundError: Could not initialize class com.datastax.spark.connector.cql.CassandraConnector

表示尚未为您的应用程序正确设置类路径。确保在启动应用程序时使用--packages 选项,如SCC Docs 中所述

关于您的其他问题

REPL 中不需要awaitTermination,因为启动流上下文后 repl 不会立即退出。该调用适用于可能没有进一步指令来阻止主线程退出的应用程序。

Start 将立即开始流式传输。

【讨论】:

  • 但我想知道为什么我能够保存下面给出的一组值 val collection = sc.parallelize(Seq(("key3", 3), ("key4", 4))) collection.saveToCassandra("test", "kv", SomeColumns("key", "value")) 在我打开的同一个 shell 中处理我的逻辑
  • 这将取决于您如何启动您的 shell 以及该操作的任务在哪里运行。该错误虽然明确告诉您类路径未在机器上正确设置。如果你有完整的跟踪,它会告诉你哪台机器。
  • 我正在使用 --packages 选项,请检查问题的“初始化 shell 的命令”部分!并且类路径已正确设置!一旦我得到那个异常,我就尝试根据那个异常检查和解决!
  • 如果它在 Repl 中运行,为什么还要构建 sbt?另外,您为什么不使用最新的 1.6.X 版本的 SCC?我对 HortonWorks 的了解还不够,无法理解他们如何改变类加载器。我至少会使用最新的 SCC 版本。 github.com/datastax/spark-cassandra-connector#169
  • 如果您查看您在回复中发布的 SSC 文档链接的 github.com/datastax/spark-cassandra-connector/blob/master/doc/… 部分,您可以理解为什么我在 REPL 上运行时使用 build.sbit!我按照相同的来源了解了如何在 Cassandra 中保存的过程!但是当我未能按照要求修改代码时,我得到了!这就是我寻求帮助的地方!
【解决方案2】:

与上下文相关的一行或两行代码导致了这里的问题!

我在浏览上下文主题时找到了解决方案!

在这里我运行了多个上下文,但它们彼此独立。

我已经用下面的命令初始化了 shell:

/usr/hdp/2.6.0.3-8/spark/bin/spark-shell --packages datastax:spark-cassandra-connector:1.6.3-s_2.10 --conf spark.cassandra.connection.host=127.0.0.1 –jars spark-streaming-kafka-assembly_2.10-1.6.3.jar

因此,当 shell 启动时,会初始化具有 Datastax 连接器属性的 spark 上下文。

后来我创建了一些配置,并使用这些配置创建了火花流上下文。使用这个上下文,我创建了 kafkaStream。这个kafkaStream只有SSC的属性,没有SC的属性,所以这里提出了存储到cassandra的问题。

我已经尝试在下面解决它并成功了!


val sc = new SparkContext(new SparkConf().setAppName("Spark-Kafka-Streaming").setMaster("local[*]").set("spark.cassandra.connection.host", "127.0.0.1"))
val ssc = new StreamingContext(sc, Seconds(10))

感谢所有站出来支持的人! 让我知道是否有更好的方法来实现它!

【讨论】:

    【解决方案3】:

    一个非常简单的方法是将流转换为 foreachRDD API 的数据帧,将 RDD 转换为 DataFrame 并使用 SparkSQL-Cassandra Datasource API 保存到 cassandra。下面是一个简单的代码 sn-p,我将 Twitter 推文保存到 Cassandra 表中

    stream.foreachRDD(rdd => {
      if (rdd.count() > 0) {
        val data = rdd.filter(status => status.getLang.equals("en")).map(status => TweetsClass(status.getId,
          status.getCreatedAt.toGMTString(),
          status.getUser.getLocation,
          status.getText)).toDF()
        //Save the data to Cassandra
        data.write.
          format("org.apache.spark.sql.cassandra").
          options(Map(
            "table" -> "sentiment_tweets",
            "keyspace" -> "My Keyspace",
            "cluster" -> "My Cluster")).mode(SaveMode.Append).save()
    
      }
    })
    

    【讨论】:

    • 请帮助我为通过流获取的数据格式制作一个简单的示例代码!请检查问题的“使用上述 Spark Streaming Context aka SSC 创建 Kafka 流”部分中的数据格式!
    • 嗨,Ishan Kumar,请查看下面的答案!我找到了解决方案并在描述的答案中解决了问题。感谢您前来提供帮助。
    猜你喜欢
    • 2016-02-07
    • 2019-06-01
    • 2019-04-29
    • 1970-01-01
    • 2016-05-01
    • 2020-04-20
    • 2022-01-23
    • 2019-07-07
    • 2019-08-13
    相关资源
    最近更新 更多