【问题标题】:How to load all records from kafka topic using spark in batch mode如何在批处理模式下使用 spark 从 kafka 主题加载所有记录
【发布时间】:2019-11-04 05:06:56
【问题描述】:

我想使用 spark 加载来自 kafka 主题的所有记录,但我看到的所有示例都使用 spark-streaming。我怎样才能只加载一次 fwom kafka 的消息?

【问题讨论】:

  • 您能否添加一个这种流式传输行为的示例和一些伪代码来说明您希望它如何工作?这表明您已经努力自己寻找解决方案,并阻止人们认为您希望人们为您编写代码。
  • 不需要,还没有收到正确答案。

标签: apache-spark apache-kafka apache-spark-sql kafka-consumer-api


【解决方案1】:

具体步骤列在in the official documentation,例如:

val df = spark
  .read
  .format("kafka")
  .option("kafka.bootstrap.servers", "host1:port1,host2:port2")
  .option("subscribePattern", "topic.*")
  .option("startingOffsets", "earliest")
  .option("endingOffsets", "latest")
  .load()

但是,如果源是连续流,则“所有记录”的定义相当不明确,因为结果取决于执行查询的时间点。

此外,您应该记住,并行性受到 Kafka 主题分区的限制,因此您必须小心不要使集群不堪重负。

【讨论】:

  • 注意:这里返回的数据只是二进制,还需要解析
猜你喜欢
  • 1970-01-01
  • 2016-10-27
  • 2021-09-20
  • 2021-09-03
  • 1970-01-01
  • 2019-07-02
  • 1970-01-01
  • 2019-07-26
  • 1970-01-01
相关资源
最近更新 更多