【问题标题】:How to fetch Kafka Stream and print it in Spark Shell?如何获取 Kafka Stream 并在 Spark Shell 中打印?
【发布时间】: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


【解决方案1】:
WARN AbstractLifeCycle: FAILED SelectChannelConnector@0.0.0.0:4040: java.net.Bind   
java.net.BindException: Address already in use

表示需要的端口已经在使用中。通常,Spark-thrifteserver 使用端口 4040。所以尝试使用 spark/sbin 文件夹中的 stop-thriftserver.sh 来停止 thriftserver。或者检查还有谁使用这个端口并释放它。

【讨论】:

  • 试过你说的那个!在启动 spark shell 之前!没有发生任何问题!但是,一旦我开始使用 spark shell 并且必须再创建一个 spark 上下文,就会再次引发相同的异常!我试图杀死那个进程,然后执行关闭!因为与之关联的进程是 java !由于它们只是警告,我认为它不会对执行产生太大影响!但主要问题为什么提取的流无法在屏幕上打印!
  • 不鼓励使用多个 Spark 上下文,这真是个坏主意。您可以使用 sc.stop() 方法停止现有的 SparkContext,然后创建一个新的。此外,您可以使用现有的 SparkContext 使用 getOrCreate(conf)
【解决方案2】:

我能够通过执行以下操作来修复它:

  1. 注意区分大小写,因为 Scala 是区分大小写的语言。 在下面的代码部分使用 map() 而不是 Map()

val kafkaStream = KafkaUtils.createStream(ssc, "localhost:2181","spark-streaming-consumer-group", map("spark-topic" -&gt; 5)) 错误做法!这就是 spark 无法获取您的流的原因!

val kafkaStream = KafkaUtils.createStream(ssc, "localhost:2181","spark-streaming-consumer-group", Map("spark-topic" -&gt; 5)) 正确练习!

  1. 检查生产者是否正在流式传输到地图功能中提到的 Kafka 主题! Spark 无法从提到的主题中获取您的流,无论是当 Kafka 生产者流数据到该主题时停止或流中的数据完成时,spark 开始从数组缓冲区中删除 RDD 并显示如下消息!

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                                                                                                                               
-------------------------------------------

  1. 关注 @YehorKrivokon 和 @VinodChandak 的 cmets 和响应,以避免面临警告!

【讨论】:

    【解决方案3】:

    有点晚了,但这可以帮助其他人。 Spark shell 已经实例化了 SparkContext,它以 sc 的形式提供。因此,要创建 StreamingContext,只需将现有的 sc 作为参数传递。希望这会有所帮助!

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2016-11-25
      • 1970-01-01
      • 1970-01-01
      • 2018-03-13
      • 1970-01-01
      • 2017-03-05
      • 2017-07-03
      • 2016-05-27
      相关资源
      最近更新 更多