【问题标题】:Does DStream's RDD pull entire data created for the batch interval at one shot?DStream 的 RDD 是否一次性提取为批处理间隔创建的全部数据?
【发布时间】:2017-03-27 00:48:53
【问题描述】:
我已经完成了this stackoverflow 问题,根据答案它创建了一个DStream,批处理间隔只有一个RDD。
例如:
我的批处理间隔是 1 分钟,Spark Streaming 作业正在使用来自 Kafka 主题的数据。
我的问题是,DStream 中可用的 RDD 是否提取/包含最后一分钟的全部数据?我们需要设置任何标准或选项来提取最后一分钟创建的所有数据吗?
如果我有一个包含 3 个分区的 Kafka 主题,并且所有 3 个分区都包含最后一分钟的数据,那么 DStream 是否会提取/包含所有 Kafka 主题分区中最后一分钟创建的所有数据?
更新:
在哪种情况下 DStream 包含多个 RDD?
【问题讨论】:
标签:
apache-spark
apache-kafka
spark-streaming
dstream
【解决方案1】:
Spark Streaming DStream 正在使用已分区的 Kafka 主题中的数据,例如 3 个不同 Kafka 代理上的 3 个分区。
DStream 中可用的 RDD 是否提取/包含最后一分钟的全部数据?
不完全是。 RDD only 描述了从提交任务执行时读取数据的偏移量。就像 Spark 中的其他 RDD 一样,它们仅 (?) 描述提交任务时要做什么以及在哪里找到要处理的数据。
但是,如果您以更宽松的方式使用“拉取/包含”来表示在某些时候将处理记录(来自给定偏移量的分区),是的,您是对的,整分钟是映射到偏移量,而偏移量又映射到 Kafka 移交给处理的记录。
在所有 Kafka 主题分区中?
是的。处理它的不一定是 Kafka 的 Spark Streaming / DStream / RDD。 DStream 的 RDD 从上次查询到现在的每个偏移量的主题及其分区请求记录。
对于 Kafka,Spark Streaming 的分钟可能略有不同,因为 DStream 的 RDD 包含偏移记录而不是每次记录。
什么情况下 DStream 包含多个 RDD?
从不。
【解决方案2】:
我建议阅读Spark documentation 中有关DStream 抽象的更多信息。
Discretized Stream 或 DStream 是 Spark Streaming 提供的基本抽象。它代表连续的数据流[...]。在内部,DStream 由一系列连续的 RDD 表示。
我要补充一点——不要忘记 RDD 本身是另一个抽象层,因此它可以分成更小的块并分布在整个集群中。
考虑您的问题:
- 是的,在每个批处理间隔触发后,都会有一个具有一个 RDD 的作业。而这个RDD包含前一分钟的所有数据。
- 如果您的作业使用具有更多分区的 Kafka 流,则所有分区都将并行使用。所以结果是所有分区的数据都在后续的RDD中处理。
【解决方案3】:
一个被忽略的重要事情是 Kafka 有多个 Spark Streaming 实现。
一种是基于接收器的方法,它在选定的 Worker 节点上设置接收器并读取数据、缓冲数据然后分发。
另一种是receiver-less 方法,这是完全不同的。它仅在运行驱动程序的节点中消耗 offsets,然后当它分配任务时,它会向每个执行程序发送一系列偏移量以供读取和处理。这样,就没有缓冲(因此,没有接收器),并且每个偏移量都由运行在工作器上的互斥执行器进程消耗。
DStream 拉取/包含所有 Kafka 主题分区中最后一分钟创建的所有数据?
在这两种方法中,它都会。每隔一分钟,它会尝试从 Kafka 读取数据并将其传播到集群中进行处理。
在这种情况下,DStream 包含多个 RDD
正如其他人所说,它永远不会。在给定的时间间隔内,只有一个 RDD 在 DStream 内流动。