【问题标题】:Spark Structured Streaming getting messages for last Kafka partitionSpark Structured Streaming 获取最后一个 Kafka 分区的消息
【发布时间】:2019-04-27 00:36:47
【问题描述】:

我正在使用 Spark Structured Streaming 从 Kafka 主题中读取数据。

无需任何分区,Spark Structired Streaming 消费者就可以读取数据。

但是当我向主题添加分区时,客户端仅显示来自最后一个分区的消息。 IE。如果主题中有 4 个分区,并且我在主题中推送 1、2、3、4 之类的数字,则客户端只打印 4 个而不是其他值。

我正在使用来自 Spark Structured Streaming 网站的最新示例和二进制文件。

    DataFrame<Row> df = spark
 .readStream()
 .format("kafka") 
.option("kafka.bootstrap.servers", "host1:port1,host2:port2") 
.option("subscribe", "topic1") 
.load()

我错过了什么吗?

【问题讨论】:

  • 如何推送消息?包括钥匙?您怎么知道这些消息甚至会发送到其他分区?此外,在您重新启动应用程序之前,Spark 不会自动拾取新分区
  • 我正在使用 Kafka 控制台生产者在主题中手动推送数据。
  • 当然,但是您使用的是--parse-keys=true吗?如果没有,您如何检查您的消息进入哪些分区?
  • 我无法检查任何特定分区。我有 4 个分区。如果我向主题发送 4 条消息,则 spark 消费者只能打印第 4 条消息。消费者仅在 4 个表中打印消息,即第 4、8、12 条消息。
  • 您可以使用Kafka的GetOffsetShell列出每个分区的最新偏移量。这会告诉你消息是否被发送到任何/所有分区......否则,如果你只有一个 Spark 执行器,那么它只会从一个 Kafka 分区消耗,所以你需要有更多

标签: java apache-spark apache-kafka spark-structured-streaming


【解决方案1】:

通过将 kafka-clients-0.10.1.1.jar 更改为 kafka-clients-0.10.0.1.jar 解决了问题。

在这里找到参考Spark Structured Stream get messages from only one partition of Kafka

【讨论】:

    猜你喜欢
    • 2017-07-03
    • 2020-12-30
    • 1970-01-01
    • 1970-01-01
    • 2021-05-05
    • 1970-01-01
    • 1970-01-01
    • 2023-03-08
    • 2018-09-26
    相关资源
    最近更新 更多