【发布时间】: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 中读取两次数据(基于我在问题中指定的用例)
标签: apache-spark spark-streaming kafka-consumer-api