【发布时间】:2015-11-17 00:50:07
【问题描述】:
我正在尝试了解 Spark Streaming 中 Spark DStream 上的转换。
我知道与地图相比,变换是最高级的,但是 有人可以给我一些可以区分变换和映射的实时示例或清晰示例吗?
【问题讨论】:
标签: apache-spark spark-streaming
我正在尝试了解 Spark Streaming 中 Spark DStream 上的转换。
我知道与地图相比,变换是最高级的,但是 有人可以给我一些可以区分变换和映射的实时示例或清晰示例吗?
【问题讨论】:
标签: apache-spark spark-streaming
Spark 流式传输中的transform 函数允许使用 Apache Spark 对流式底层RDDs 的任何转换。 map 用于元素到元素的转换,可以使用transform 实现。本质上,map 作用于DStream 的元素,transform 允许您使用 DStream 的RDDs。您可能会发现http://spark.apache.org/docs/latest/streaming-programming-guide.html#transformations-on-dstreams 很有用。
【讨论】:
map 是基本变换,transform 是 RDD 变换
map(func) : 通过传递源的每个元素返回一个新的 DStream 通过函数 func 进行 DStream。
这是一个演示 DStream 上的映射操作和变换操作的示例
val conf = new SparkConf().setMaster("local[*]").setAppName("StreamingTransformExample")
val ssc = new StreamingContext(conf, Seconds(5))
val rdd1 = ssc.sparkContext.parallelize(Array(1,2,3))
val rdd2 = ssc.sparkContext.parallelize(Array(4,5,6))
val rddQueue = new Queue[RDD[Int]]
rddQueue.enqueue(rdd1)
rddQueue.enqueue(rdd2)
val numsDStream = ssc.queueStream(rddQueue, true)
val plusOneDStream = numsDStream.map(x => x+1)
plusOneDStream.print()
map 操作将 DStream 中所有 RDD 中的每个元素加 1,输出如下所示
-------------------------------------------
Time: 1501135220000 ms
-------------------------------------------
2
3
4
-------------------------------------------
Time: 1501135225000 ms
-------------------------------------------
5
6
7
-------------------------------------------
transform(func) : 通过应用 RDD-to-RDD 函数返回一个新的 DStream 到源 DStream 的每个 RDD。这可以用来做任意 DStream 上的 RDD 操作。
val commonRdd = ssc.sparkContext.parallelize(Array(0))
val combinedDStream = numsDStream.transform(rdd=>(rdd.union(commonRdd)))
combinedDStream.print()
transform 允许对 DStream 中的 RDD 执行连接、联合等 RDD 操作,此处给出的示例代码将产生如下输出
-------------------------------------------
Time: 1501135490000 ms
-------------------------------------------
1
2
3
0
-------------------------------------------
Time: 1501135495000 ms
-------------------------------------------
4
5
6
0
-------------------------------------------
Time: 1501135500000 ms
-------------------------------------------
0
-------------------------------------------
Time: 1501135505000 ms
-------------------------------------------
0
-------------------------------------------
这里包含元素 0 的 commonRdd 与 DStream 中的所有底层 RDD 执行联合操作。
【讨论】:
DStream 有多个 RDD,因为每个批次间隔都是不同的 RDD。
因此,通过使用 transform(),您有机会在
整个 DStream。
来自 Spark Docs 的示例: http://spark.apache.org/docs/latest/streaming-programming-guide.html#transform-operation
【讨论】:
Spark Streaming 中的转换函数允许您对 Stream 中的底层 RDD 执行任何转换。例如,您可以使用 Transform 在流中加入两个 RDD,其中一个 RDD 是由文本文件或并行集合制成的一些 RDD,而其他 RDD 来自文本文件/套接字等流。
Map 在特定批次中作用于 RDD 的每个元素,并在应用传递给 Map 的函数后生成 RDD。
【讨论】:
示例 1)
男人们排队进入房间,换衣服,然后娶他们选择的女人。
1) 换衣服是贴图操作(就是在属性中变换自己)
2) 娶女是对你的合并/过滤操作,但受到其他人的影响,我们可以称之为真正的变换操作。
示例 2) 学生进入大学,很少人参加了 2 节课,很少有人参加了 4 节课,以此类推。
1) 听课是地图操作,一般是学生在做的事情。
2) 但要确定讲师教给他们的内容取决于讲师 RDD 数据,即他的日程安排。
假设转换操作是您要过滤或验证的维度或静态表,以便为您识别正确的数据,而删除可能是垃圾。
【讨论】:
如果我有来自 0-1 秒 "Hello How" 和接下来 1-2 秒 "Are You" 的数据。然后在 map 和 reduce by key 示例的情况下,如上所示,将为第一批生成输出 (hello,1) 和 (How,1),为下一批生成 (are,1) 和 (you,1)。但是同样对于下一个使用“变换函数”的示例,输出会有什么不同
【讨论】: