【问题标题】:Aggregate data from different micro batches in Spark streaming在 Spark 流中聚合来自不同微批次的数据
【发布时间】:2017-09-16 16:18:44
【问题描述】:

我正在尝试每分钟使用 Spark 流式传输(从 Kafka 读取)聚合并查找一些指标。我能够汇总那一分钟的数据。如何确保我可以拥有当天的存储桶并总结当天所有分钟的所有汇总值?

我有一个数据框,我正在做类似的事情。

sampleDF = spark.sql("select userId,sum(likes) as total from likes_dataset group by userId order by userId")

【问题讨论】:

    标签: apache-spark pyspark apache-spark-sql spark-streaming spark-dataframe


    【解决方案1】:

    您可以使用结构化流式编程中的“Watermarking”功能

    示例代码

    import spark.implicits._
    
    val words = ... // streaming DataFrame of schema { timestamp: Timestamp, word: String }
    
        val windowedCounts = words
            .withWatermark("timestamp", "10 minutes")
            .groupBy(
                window($"timestamp", "10 minutes", "5 minutes"),
                $"word")
            .count()
    

    【讨论】:

    • 感谢您的回答。我试过这样做。 Spark 不保留之前的微批量值。如果我有 60 秒的微批处理间隔并且如果我尝试创建 10 分钟的窗口,则 12:01:00 的值不会与 12:02:00 的值聚合。对于 12:02:00,它只查找最近收到的数据的聚合。如何存储所有 10 分钟的汇总数据?
    • 我有从 Kafka 流中获取数据的主要功能。并且,它为每个 RDD 调用一个函数。在这个函数中,我确实聚合了这些值。但是对于每个 RDD,聚合值都会被重置。以前的聚合值不会保留。我不知道如何为这个 spark 会话生命定义一个全局聚合数据框并组合所有聚合数据。有人可以帮忙吗?
    【解决方案2】:

    我知道发生了什么。我开始了解 Spark 中的状态流,这对我有帮助。

    我所要做的就是,

    running_counts = countStream.updateStateByKey(updateTotalCount, initialRDD=initialStateRDD) 
    

    我不得不写这个 updateTotalCount 函数来说明如何将旧聚合数据与微批处理的新聚合数据合并。就我而言,更新函数如下所示:

    def updateTotalCount(currentCount, countState):
        if countState is None:
           countState = 0
        return sum(currentCount) + countState
    

    【讨论】:

    • 感谢您的回答。你能在countState 上解释一下吗,它有什么?还有你的countStream 有一个RDD 流?你怎么得到initialRDDinitialStateRDD
    猜你喜欢
    • 2023-03-23
    • 1970-01-01
    • 1970-01-01
    • 2021-08-21
    • 1970-01-01
    • 1970-01-01
    • 2020-07-22
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多