Guo*_*Guo 3 apache-spark spark-streaming
在火花流中,数据根据批处理间隔进行处理.如果我设置一个5秒的批处理间隔(val ssc = new StreamingContext(sc, Seconds(5))):
1s~5s is first batch of data
6s~10s is second batch of data
10s~15s is third batch of data
……
Run Code Online (Sandbox Code Playgroud)
是否有变量来识别火花流中的每个批次数据?如果有这样的变量:
var batchID = 0
我可以获得batchID识别哪一批数据的价值,或者我可以通过batchID过滤数据,如:window(……).filter(_.batchId == 1).
或者有没有办法区分每批数据?
您可以使用foreachRDD哪种类型(rdd: RDD[T], time: Time) => Unit.时间是RDD数据流中的标记,表示在两个连续批次上的两次连续调用中,时间参数将相差一个批次间隔持续时间.
您可以在foreachRDD此处找到API :https:
//spark.apache.org/docs/latest/api/scala/index.html#org.apache.spark.streaming.dstream.DStream
如果您需要RDD为特定的时间间隔选择一些s,您只需使用该slice功能,该功能也在上面的链接中指定.
| 归档时间: |
|
| 查看次数: |
1872 次 |
| 最近记录: |