【发布时间】:2021-11-16 16:42:44
【问题描述】:
Kafka enable.auto.commit 设置为 false,Spark 版本为 2.4
-
如果使用最新的偏移量,我们是否需要手动查找最后的偏移量详细信息并在 Spark 应用程序的 .CreateDirectStream() 中提及?还是会自动采用最新的偏移量?无论如何,我们是否需要手动查找最后的偏移量详细信息。
-
使用
SparkSession.readstrem.format(kafka)....和KafkaUtils.createDirectStream()有什么区别吗? -
使用最早偏移选项时,会自动考虑偏移吗?
【问题讨论】:
-
kafkadetails到底是什么?那个物体是什么?来自哪个图书馆? -
是KafkaUtils,我会在描述中更正
-
感谢您的更新,从您所说的我可以看出您使用的是 Spark 2.2 或更低版本,因为 KafkaUtils 在该版本之后不可用。另一个问题是这里的
spark是SparkContext还是SparkSession,因为它们提供了不同的API。 -
我使用的是 spark 2.4
-
例如:val stream = KafkaUtils.createDirectStream[String, String](streamingContext, PreferConsistent, Subscribe[String, String](topics, kafkaParams))
标签: scala apache-spark apache-kafka