【发布时间】:2019-04-03 09:46:50
【问题描述】:
如何通过指定开始和结束偏移量将kafka主题中的数据读取到RDD?
KafkaUtils.createRDD is 是实验性的,API 相当令人不快(它返回一个臃肿的 Java ConsumerRecord 类,它甚至不能序列化并将它放在 KafkaRDD 中,它覆盖了很多的方法(如持久化)只是抛出一个异常。
我想要的是这样一个简单的 API:
case class Message(key: String,
value: String,
offset: Long,
timestamp: Long)
def readKafka(topic: String, offsetsByPartition: Map[Int, (Long, Long)])
(config: KafkaConfig, sc: SparkContext): RDD[Message]
或者类似key: Array[Byte]和value: Array[Byte]的东西
【问题讨论】:
标签: scala apache-spark apache-kafka spark-structured-streaming