Bry*_*yan 5 apache-kafka apache-spark spark-streaming apache-zookeeper
我试图通过火花流来读取来自Kafka的旧消息.但是,我只能在实时发送消息时检索消息(即,如果我填充新消息,而我的spark程序正在运行 - 那么我会收到这些消息).
我正在更改我的groupID和consumerID,以确保zookeeper不仅不会发出它知道我的程序以前见过的消息.
假设spark将zookeeper中的偏移量视为-1,那么它是否应该读取队列中的所有旧消息?我只是误解了kafka队列的使用方式吗?我很新兴火花和卡夫卡,所以我不能排除我只是误解了一些东西.
package com.kibblesandbits
import org.apache.spark.SparkContext
import org.apache.spark.streaming.{Seconds, StreamingContext}
import org.apache.spark.streaming.kafka.KafkaUtils
import net.liftweb.json._
object KafkaStreamingTest {
val cfg = new ConfigLoader().load
val zookeeperHost = cfg.zookeeper.host
val zookeeperPort = cfg.zookeeper.port
val zookeeper_kafka_chroot = cfg.zookeeper.kafka_chroot
implicit val formats = DefaultFormats
def parser(json: String): String = {
return json
}
def main(args : Array[String]) {
val zkQuorum = "test-spark02:9092"
val group = "myGroup99"
val topic = Map("testtopic" -> 1)
val sparkContext = new SparkContext("local[3]", "KafkaConsumer1_New")
val ssc = new StreamingContext(sparkContext, Seconds(3))
val json_stream = KafkaUtils.createStream(ssc, zkQuorum, group, topic)
var gp = json_stream.map(_._2).map(parser)
gp.saveAsTextFiles("/tmp/sparkstreaming/mytest", "json")
ssc.start()
}
Run Code Online (Sandbox Code Playgroud)
运行时,我将看到以下消息.所以我相信它不仅没有看到消息,因为设置了偏移量.
14/12/05 13:34:08 INFO ConsumerFetcherManager:[ConsumerFetcherManager-1417808045047]为分区添加了fetcher ArrayBuffer([[testtopic,0],initOffset -1到代理id:1,host:test-spark02.vpc,port: 9092],[[testtopic,1],initOffset -1到代理ID:1,主机:test-spark02.vpc,端口:9092],[[testtopic,2],initOffset -1到代理ID:1,主机: test-spark02.vpc,port:9092],[[testtopic,3],initOffset -1到broker id:1,host:test-spark02.vpc,port:9092],[[testtopic,4],initOffset -1经纪人ID:1,主持人:test-spark02.vpc,port:9092])
然后,如果我填充1000条新消息 - 我可以看到我的临时目录中保存的那1000条消息.但我不知道如何阅读现有的消息,这些消息应该在(此时)成千上万.
使用替代工厂方法KafkaUtils,可以为Kafka使用者提供配置:
def createStream[K: ClassTag, V: ClassTag, U <: Decoder[_]: ClassTag, T <: Decoder[_]: ClassTag](
ssc: StreamingContext,
kafkaParams: Map[String, String],
topics: Map[String, Int],
storageLevel: StorageLevel
): ReceiverInputDStream[(K, V)]
Run Code Online (Sandbox Code Playgroud)
然后使用kafka配置构建一个映射,并将参数'kafka.auto.offset.reset'设置为'smallest':
val kafkaParams = Map[String, String](
"zookeeper.connect" -> zkQuorum, "group.id" -> groupId,
"zookeeper.connection.timeout.ms" -> "10000",
"kafka.auto.offset.reset" -> "smallest"
)
Run Code Online (Sandbox Code Playgroud)
将该配置提供给上面的工厂方法."kafka.auto.offset.reset" - >"smallest"告诉消费者从主题中的最小偏移量开始.
| 归档时间: |
|
| 查看次数: |
3397 次 |
| 最近记录: |