【问题标题】:Is there a variable to identify each batch data in spark streaming?是否有一个变量来识别火花流中的每个批次数据?
【发布时间】:2016-02-02 17:34:30
【问题描述】:

在 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
……

是否有一个变量来识别火花流中的每个批次数据?如果有这样的变量:

var batchID = 0

我可以获取batchID 的值来识别哪一批数据,或者我可以按batchID 过滤数据,例如:window(……).filter(_.batchId == 1)

或者有什么方法可以区分每批数据?

【问题讨论】:

    标签: apache-spark spark-streaming


    【解决方案1】:

    您可以使用类型为(rdd: RDD[T], time: Time) => UnitforeachRDD。时间是数据流中RDD 的标记,这意味着在两个连续批次的两次连续调用中,时间参数将相差一个批次间隔持续时间。

    您可以在此处找到foreachRDD 的 API: https://spark.apache.org/docs/latest/api/scala/index.html#org.apache.spark.streaming.dstream.DStream

    如果您需要为特定的时间间隔选择一些RDDs,您可以简单地使用slice函数,该函数也在上面的链接中指定。

    【讨论】:

    • 感谢您的帮助!我还有一个问题:foreachRDD 返回Unit 而不是DStream,我无法将结果打印到控制台。使用foreachRDD函数有什么方法可以轻松调试?
    • 您可以在foreachRDD内拨打println
    • 当我使用println like:.foreachRDD(println(_)),结果是MapPartitionsRDD[1] at map at SparkStreaming.scala:19,我该怎么办?
    • 开始你的 Spark Streaming 上下文了吗? (在 StreamingContext 的 API 中查找方法 start、stop、awaitTermination)
    • 我使用.foreachRDD(_.collect.foreach(println)) 并且工作得很好。谢谢!
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-10-02
    • 1970-01-01
    • 1970-01-01
    • 2015-01-30
    • 1970-01-01
    • 2022-01-19
    • 1970-01-01
    相关资源
    最近更新 更多