Skip to content

Instantly share code, notes, and snippets.

@rssanders3
Created January 25, 2019 15:59
Show Gist options
  • Select an option

  • Save rssanders3/22a55cf0d5fb22db2f5e581f15f130e0 to your computer and use it in GitHub Desktop.

Select an option

Save rssanders3/22a55cf0d5fb22db2f5e581f15f130e0 to your computer and use it in GitHub Desktop.
val storedOffsets: Option[mutable.Map[TopicPartition, Long]] = loadOffsets(spark, kuduContext)
val kafkaDStream = storedOffsets match {
case None =>
LOGGER.info("storedOffsets was None")
kafkaParams += ("auto.offset.reset" -> "latest")
KafkaUtils.createDirectStream[String, Array[Byte]]
(ssc, PreferConsistent, ConsumerStrategies.Subscribe[String, Array[Byte]]
(topicsSet, kafkaParams)
)
case Some(fromOffsets) =>
LOGGER.info("storedOffsets was Some(" + fromOffsets + ")")
kafkaParams += ("auto.offset.reset" -> "none")
KafkaUtils.createDirectStream[String, Array[Byte]]
(ssc, PreferConsistent, ConsumerStrategies.Assign[String, Array[Byte]]
(fromOffsets.keys.toList, kafkaParams, fromOffsets)
)
}
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment