【发布时间】:2015-07-14 22:24:56
【问题描述】:
Spark 流式处理微批量的数据。
使用 RDD 并行处理每个区间数据,每个区间之间不共享任何数据。
但我的用例需要在区间之间共享数据。
考虑Network WordCount 示例,该示例生成在该时间间隔内接收到的所有单词的计数。
我将如何产生以下字数?
单词“hadoop”和“spark”与前一个间隔计数的相对计数
所有其他单词的正常字数。
注意:UpdateStateByKey 会进行有状态处理,但这会将函数应用于每条记录而不是特定记录。
所以,UpdateStateByKey 不符合这个要求。
更新:
考虑以下示例
间隔 1
输入:
Sample Input with Hadoop and Spark on Hadoop
输出:
hadoop 2
sample 1
input 1
with 1
and 1
spark 1
on 1
间隔 2
输入:
Another Sample Input with Hadoop and Spark on Hadoop and another hadoop another spark spark
输出:
another 3
hadoop 1
spark 2
and 2
sample 1
input 1
with 1
on 1
说明:
第一个区间给出所有单词的正常字数。
在第二个时间间隔内,hadoop 发生了 3 次,但输出应该是 1 (3-2)
火花发生 3 次,但输出应为 2 (3-1)
对于所有其他单词,它应该给出正常的字数。
因此,在处理 2nd Interval 数据时,它应该具有 hadoop 和 spark 的 1st interval 字数
这是一个带有插图的简单示例。
在实际用例中,需要数据共享的字段是RDD元素(RDD)的一部分,需要跟踪的值非常多。
即在这个例子中像 hadoop 和 spark 关键字一样有近 100k 个关键字被跟踪。
Apache Storm 中的类似用例:
【问题讨论】:
-
请澄清““hadoop”和“spark”这两个词的相对计数与前一个间隔计数。不要犹豫,将您的定义形式化,引入变量和公式。你也可以举个例子。
-
过滤掉不需要的并updateStateByKey?
标签: apache-spark spark-streaming