【发布时间】:2018-04-04 21:01:19
【问题描述】:
首先我通过以下方式在文件夹中构建了一个 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"
后来在我的 "build.sbt" 所在的同一个文件夹中,我按以下方式启动了 spark 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
这些是 spark shell 启动时显示的警告:
WARN AbstractLifeCycle: FAILED SelectChannelConnector@0.0.0.0:4040: java.net.Bind
java.net.BindException: Address already in use
WARN AbstractLifeCycle: FAILED org.spark-project.jetty.server.Server@75bf9e67:
java.net.BindException: Address already in use
在 Spark shell 中导入以下包
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 shell 中按以下方式创建配置:
val conf = new SparkConf().setMaster("local[*]").setAppName("KafkaReceiver").set("spark.driver.allowMultipleContexts", "true").setMaster("local");
在创建分配它的配置后,通过以下方式创建一个新的火花流上下文:
val ssc = new StreamingContext(conf, Seconds(10))
在创建火花流上下文期间,上面显示的一些警告与其他警告一起再次出现,如下所示
WARN AbstractLifeCycle: FAILED SelectChannelConnector@0.0.0.0:4040: java.net.Bind
java.net.BindException: Address already in use
.
.
.
WARN AbstractLifeCycle: FAILED org.spark-project.jetty.server.Server@75bf9e67:
java.net.BindException: Address already in use
.
.
.
WARN SparkContext: Multiple running SparkContexts detected in the same JVM!
org.apache.spark.SparkException: Only one SparkContext may be running in this JVM (see SPARK-2243). To ignore this error, set spark.driver.allowMulti
pleContexts = true. The currently running SparkContext was created at:
org.apache.spark.SparkContext.<init>(SparkContext.scala:82)
org.apache.spark.repl.SparkILoop.createSparkContext(SparkILoop.scala:1017)
.
.
.
WARN StreamingContext: spark.master should be set as local[n], n > 1 in local mode if you have receivers to get data, otherwise Spa
rk jobs will not get resources to process the received data.
ssc: org.apache.spark.streaming.StreamingContext = org.apache.spark.streaming.StreamingContext@616f1c2e
然后使用创建的火花流上下文以下面的方式创建了一个 kafkaStream
val kafkaStream = KafkaUtils.createStream(ssc, "localhost:2181","spark-streaming-consumer-group", map("spark-topic" -> 5))
然后打印流并以下面的方式启动 ssc
kafkaStream.print()
ssc.start
在shell中使用上述命令后,输出如下图所示
重复打印的输出如下所示!
17/08/18 10:01:30 INFO JobScheduler: Starting job streaming job 1503050490000 ms.0 from job set of time 1503050490000 ms
17/08/18 10:01:30 INFO JobScheduler: Finished job streaming job 1503050490000 ms.0 from job set of time 1503050490000 ms
17/08/18 10:01:30 INFO JobScheduler: Total delay: 0.003 s for time 1503050490000 ms (execution: 0.000 s)
17/08/18 10:01:30 INFO BlockRDD: Removing RDD 3 from persistence list
17/08/18 10:01:30 INFO KafkaInputDStream: Removing blocks of RDD BlockRDD[3] at createStream at <console>:39 of time 1503050490000 ms
17/08/18 10:01:30 INFO ReceivedBlockTracker: Deleting batches ArrayBuffer(1503050470000 ms)
17/08/18 10:01:30 INFO InputInfoTracker: remove old batch metadata: 1503050470000 ms
17/08/18 10:01:30 INFO BlockManager: Removing RDD 3
17/08/18 10:01:40 INFO JobScheduler: Added jobs for time 1503050500000 ms
-------------------------------------------
Time: 1503050500000 ms
-------------------------------------------
【问题讨论】:
-
你能在这里粘贴输出吗?由于访问限制,我无法打开图像。 :(
-
@VinodChandak 这些是 1000 行!我想如果我在这里粘贴会很乱!你能告诉我任何其他为你提供输出的替代方式吗!
-
Edited the question output of the finale part is added below all images @VinodChandak 希望你能通过它并尽可能帮助我!
-
我在这里没有看到任何异常。代码看起来很完美。您可以尝试添加最后一条语句。 ssc.awaitTermination() 等待终止。
-
除此之外,不要打印流本身,尝试将 kafkaStream (DStream) 持久化到文件系统。您可以为此使用 DStream.foreachRDD 方法。如果您愿意,我可以发布示例。
标签: scala shell apache-spark apache-kafka spark-streaming