【问题标题】:How to read from specific Kafka partition in Spark structured streaming如何从 Spark 结构化流中的特定 Kafka 分区中读取数据
【发布时间】:2019-02-15 06:35:53
【问题描述】:

我的 Kafka 主题有三个分区,我想知道是否可以只读取三个分区中的一个。我的消费者是 spark 结构化流应用程序。

以下是我在 spark 中现有的 kafka 设置。

  val inputDf = spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", brokers)
  .option("subscribe", topic)
  .option("startingOffsets", "latest")
  .load()

【问题讨论】:

    标签: apache-kafka kafka-consumer-api spark-structured-streaming spark-streaming-kafka


    【解决方案1】:

    这是您如何从特定分区读取数据的方法。

     val inputDf = spark.readStream
      .format("kafka")
      .option("kafka.bootstrap.servers", brokers)
      .option("assign", """{"topic":[0]}""") 
      .option("startingOffsets", "latest")
      .load()
    

    PS:从多个分区而不是 1--> """{"topic":[0,1,2..n]}"""

    【讨论】:

      【解决方案2】:

      同样,您如何写入特定分区。我试过了,它不起作用。

              someDF
                .selectExpr("key", "value")
                .writeStream
                .format("kafka")
                .option("kafka.bootstrap.servers", kafkaServers)
                .option("topic", "someTopic")
                .option("partition", partIdx)
                .start()
      

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2019-08-01
        • 2018-08-20
        • 1970-01-01
        • 2016-06-09
        • 1970-01-01
        • 2017-04-04
        • 1970-01-01
        • 2019-01-29
        相关资源
        最近更新 更多