spark-user mailing list archives

Site index · List index
Message view « Date » · « Thread »
Top « Date » · « Thread »
From Jatin Kumar <>
Subject Spark streaming from Kafka best fit
Date Tue, 01 Mar 2016 08:36:48 GMT
Hello all,

I see that there are as of today 3 ways one can read from Kafka in spark
1. KafkaUtils.createStream() (here
2. KafkaUtils.createDirectStream() (here
3. Kafka-spark-consumer (here

My spark streaming application has to read from 1 kafka topic with around
224 partitions, consuming data at around 150MB/s (~90,000 messages/sec)
which reduces to around 3MB/s (~1400 messages/sec) after filtering. After
filtering I need to maintain top 10000 URL counts. I don't really care
about exactly once semantics as I am interested in rough estimate.


sparkConf.set("spark.streaming.receiver.writeAheadLog.enable", "false")
val ssc = StreamingContext.getOrCreate(kCheckPointDir, createStreamingContext)

val topicMap = topics.split(",").map((_, numThreads.toInt)).toMap
val kafkaParams = Map[String, String](
  "" -> "kafka.server.ip:9092",
  "" -> consumer_group

val lineStreams = (1 to N).map{ _ =>
  KafkaUtils.createStream[String, String, StringDecoder, StringDecoder](
    ssc, kafkaParams, topicMap, StorageLevel.MEMORY_AND_DISK).map(_._2)

ssc.union( => {
    .filter(record => isGoodRecord(record))
    .map(record => record.url)
).window(Seconds(120), Seconds(120))  // 2 Minute window
  .countByValueAndWindow(Seconds(1800), Seconds(120), 28) // 30 Minute
moving window, 28 will probably help in parallelism
  .filter(urlCountTuple => urlCountTuple._2 > MIN_COUNT_THRESHHOLD)
  .mapPartitions(iter => {
    iter.toArray.sortWith((r1, r2) => <
0).slice(0, 1000).iterator
  }, true)
  .foreachRDD((latestRDD, rddTime) => {
      printTopFromRDD(rddTime, => (record._2,



a) I used #2 but I found that I couldn't control how many executors will be
actually fetching from Kafka. How do I keep a balance of executors which
receive data from Kafka and which process data? Do they keep changing for
every batch?

b) Now I am trying to use #1 creating multiple DStreams, filtering them and
then doing a union. I don't understand why would the number of events
processed per 120 seconds batch will change drastically. PFA the events/sec
graph while running with 1 receiver. How to debug this?

c) What will be the most suitable method to integrate with Kafka from above
3? Any recommendations for getting maximum performance, running the
streaming application reliably in production environment?

Jatin Kumar

View raw message