【发布时间】:2021-06-07 15:25:11
【问题描述】:
我想批量消费一个 Kafka 主题,我想每小时读取 Kafka 主题并读取最新的每小时数据。
val readStream = existingSparkSession
.read
.format("kafka")
.option("kafka.bootstrap.servers", hostAddress)
.option("subscribe", "kafka.raw")
.load()
但这总是读取前 20 个数据行,并且这些行从一开始就开始,所以这永远不会选择最新的数据行。
如何使用 scala 和 spark 每小时读取最新行?
【问题讨论】:
-
为什么要通过kafka做批处理?它不是在您的数据湖中的某个地方可用吗?
标签: apache-spark apache-kafka spark-streaming spark-structured-streaming