Advertisement
Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- val streams = (1 to inputTopic.partitions) map { _ =>
- KafkaUtils.createStream[Array[Byte], Array[Byte], DefaultDecoder, DefaultDecoder](
- ssc,
- kafkaParams,
- Map(inputTopic.name -> 1),
- storageLevel = StorageLevel.MEMORY_ONLY_SER // or: StorageLevel.MEMORY_AND_DISK_SER
- ).map(_._2)
- }
Advertisement
Add Comment
Please, Sign In to add comment
Advertisement