【问题标题】:Simple Spark Structured Streaming equivalent of KafkaUtils.createRDD, i.e. read kafka topic to RDD by specifying offsets?KafkaUtils.createRDD 的简单 Spark Structured Streaming 等价物,即通过指定偏移量将 kafka 主题读取到 RDD?
【发布时间】: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


    【解决方案1】:

    要从带有偏移量的 kafka 中读取,代码如下所示,参考 here

    val df = 
      spark.readStream
      .format("kafka")
      .option("kafka.bootstrap.servers", "host1:port1,host2:port2")
      .option("subscribe", "topic1,topic2")
      .option("startingOffsets", """{"topic1":{"0":23,"1":-2},"topic2":{"0":-2}}""")
      .option("endingOffsets", """{"topic1":{"0":50,"1":-1},"topic2":{"0":-1}}""")
      .load()
    

    上面将读取偏移量内可用的数据,然后您可以将列转换为字符串,并转换为您的对象Message

    val messageRDD: RDD[Message] = 
      df.select(
        col("key").cast("string"), 
        col("value").cast("string"), 
        col("offset").cast("long"),
        col("timestamp").cast("long")
      ).as[Message].rdd
    

    【讨论】:

      猜你喜欢
      • 2020-09-03
      • 2021-05-22
      • 2019-07-30
      • 2019-10-15
      • 1970-01-01
      • 2018-04-26
      • 2017-06-22
      • 2017-12-31
      • 1970-01-01
      相关资源
      最近更新 更多