【发布时间】: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