【问题标题】:How to get the current batch timestamp in Spark streaming如何在 Spark 流中获取当前批处理时间戳
【发布时间】:2016-03-29 23:42:19
【问题描述】:

如何在 Spark Streaming 中获取当前的批处理时间戳(DStream)?

我有一个 Spark 流应用程序,其中输入数据将经过多次转换。

我需要执行期间的当前时间戳来验证输入数据中的时间戳。

如果我与当前时间进行比较,则时间戳可能与每个 RDD 转换执行不同。

有什么方法可以获取时间戳,特定 Spark 流式微批处理何时开始或它属于哪个微批处理间隔?

【问题讨论】:

  • 嗨,你找到答案了吗?

标签: java apache-spark spark-streaming


【解决方案1】:
dstream.foreachRDD((rdd, time)=> {
  // time is scheduler time for the batch job.it's interval was your window/slide length.
})

【讨论】:

  • 这适用于 Scala,但不幸的是,Pyspark 似乎还没有等价物:spark.apache.org/docs/latest/api/python/…。我不得不在传递给foreachRDD() 的函数中使用datetime.datetime.now() 作为解决方法。
【解决方案2】:
dstream.transform(
    (rdd, time) => {
        rdd.map(
            (time, _)
        )
    }
).filter(...)

【讨论】:

    【解决方案3】:

    迟到的回复......但如果它对某人有帮助,时间戳可以提取为毫秒。首先定义一个使用Java API格式化的函数:

    使用 Java 7 - 样式 util.Date/DateFormat:

    def returnFormattedTime(ts: Long): String = {
        val date = new Date(ts)
        val formatter = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss")
        val formattedDate = formatter.format(date)
        formattedDate
    }
    

    或者,使用 Java 8 - 风格的 util.time:

    def returnFormattedTime(ts: Long): String = {
        val date = Instant.ofEpochMilli(ts)
        val formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss").withZone(ZoneId.systemDefault())
        val formattedDate = formatter.format(date)
        formattedDate
    }
    

    最后使用foreachRDD方法获取时间戳:

    dstreamIns.foreachRDD((rdd, time) =>
        ....
        println(s"${returnFormattedTime(time.milliseconds)}")
        ....
    )
    

    【讨论】:

    • 有 spark sql 流的例子吗?
    猜你喜欢
    • 1970-01-01
    • 2011-02-16
    • 2017-02-15
    • 1970-01-01
    • 2021-12-04
    • 1970-01-01
    • 2018-09-26
    • 1970-01-01
    • 2017-02-11
    相关资源
    最近更新 更多