【问题标题】:Caching DStream in Spark Streaming在 Spark Streaming 中缓存 DStream
【发布时间】:2016-10-07 15:49:48
【问题描述】:

我有一个从 kafka 读取数据的 Spark 流式处理, 进入 DStream。

在我的管道中,我做了两次(一个接一个):

DStream.foreachRDD(RDD 上的转换并插入目标)。

(每次我做不同的处理并将数据插入到不同的目的地)。

我想知道在我从 Kafka 读取数据之后,DStream.cache 将如何工作?有可能吗?

这个过程现在实际上是从 Kafka 读取数据两次吗?

请记住,不可能将两个 foreachRDD 合二为一(因为两条路径完全不同,那里有全状态转换 - 需要在 DStream 上应用...)

感谢您的帮助

【问题讨论】:

  • Dstream.cache 将起作用。它在第一次看到动作时缓存流。对于 DStream 中的后续操作,它使用缓存。
  • @Knight71 我还需要设置DStream.unpersist(true),和缓存RDD的时候一样,最后当不再需要DStream的时候?
  • 所有操作后Dstream数据会自动清除,由spark streaming根据transformation决定。
  • @Knight71,谢谢你的回答?如果我不放 DStream.cache,是否意味着 Spark 会从 Kafka 中读取两次数据(基于我在问题中指定的用例)
  • 另一个有用的链接:stackoverflow.com/questions/30253897/…

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


【解决方案1】:

有两种选择:

  • 使用Dstream.cache() 将底层RDD 标记为已缓存。 Spark Streaming 将负责在超时后取消持久化 RDD,由 spark.cleaner.ttl 配置控制。

  • 使用额外的foreachRDD 将cache() 和unpersist(false) 副作用操作应用于DStream 中的RDD:

例如:

val kafkaDStream = ???
val targetRDD = kafkaRDD
                       .transformation(...)
                       .transformation(...)
                       ...
// Right before the lineage fork mark the RDD as cacheable:
targetRDD.foreachRDD{rdd => rdd.cache(...)}
targetRDD.foreachRDD{do stuff 1}
targetRDD.foreachRDD{do stuff 2}
targetRDD.foreachRDD{rdd => rdd.unpersist(false)}

请注意,如果可以选择,您可以将缓存合并为 do stuff 1 的第一条语句。

我更喜欢这个选项,因为它可以让我对缓存生命周期进行细粒度控制,并让我在需要时立即清理内容,而不是依赖 ttl。

【讨论】:

  • spark.cleaner.ttl 已删除。这是什么新的属性控件?
猜你喜欢
  • 2020-06-03
  • 1970-01-01
  • 2016-06-12
  • 2018-07-02
  • 2015-11-03
  • 2019-03-18
  • 1970-01-01
  • 2016-12-22
  • 2014-12-21
相关资源
最近更新 更多