【问题标题】:Multiple aggregations in Spark Structured StreamingSpark Structured Streaming 中的多个聚合
【发布时间】:2017-04-22 00:21:58
【问题描述】:

我想在 Spark Structured Streaming 中进行多个聚合。

类似这样的:

  • 读取输入文件流(从文件夹中)
  • 执行聚合 1(带有一些转换)
  • 执行聚合 2(以及更多转换)

当我在结构化流中运行它时,它给了我一个错误“流数据帧/数据集不支持多个流聚合”。

有没有办法在结构化流中进行这样的多重聚合?

【问题讨论】:

  • 你试过使用较低级别的DStream抽象吗?
  • 我希望使用结构化流(数据集/数据帧)。你能指出一些使用 DStream 完成类似操作的示例吗?
  • 这个问题有什么解决办法吗?请在此处提供..相同的问题

标签: apache-spark apache-spark-sql spark-structured-streaming


【解决方案1】:

TLDR - 不支持;在某些情况下,解决方法是可能的。

加长版-

  1. (一个黑客)

在某些情况下,解决方法是可能的,例如,如果您希望在低基数列的流式查询中拥有多个 count(distinct),那么 approx_count_distinct 很容易实际返回 exact 通过将 rsd 参数设置得足够低来获得不同元素的数量(这是 approx_count_distinct 的第二个可选参数,默认情况下为 0.05)。

这里如何定义“低基数”?对于具有超过 1000 个唯一值的列,我不建议使用这种方法。

因此,在您的流式查询中,您可以执行以下操作 -

(spark.readStream....
      .groupBy("site_id")
      .agg(approx_count_distinct("domain", 0.001).alias("distinct_domains")
         , approx_count_distinct("country", 0.001).alias("distinct_countries")
         , approx_count_distinct("language", 0.001).alias("distinct_languages")
      )
  )

这是它确实有效的证据:

请注意 count(distinct) 和 count_approx_distinct 给出相同的结果! 以下是关于rsd 参数count_approx_distinct 的一些指导:

  • 对于具有 0.02 的 100 个不同值 rsd 的列是必需的;
  • 对于具有 0.001 的 1000 个不同值 rsd 的列是必需的。

PS。另请注意,我必须在具有 10k 个不同值的列上注释掉实验,因为我没有足够的耐心来完成它。这就是为什么我提到你不应该对具有超过 1k 个不同值的列使用这个 hack。对于 approx_count_distinct 来匹配超过 1k 个不同值的精确计数(不同),对于 HyperLogLogPlusPlus algorithm 的设计目标而言,rsd 太低了(该算法落后于 approx_count_distinct 实现)。

  1. (很好,但更多涉及的方式)

正如其他人提到的,您可以使用 Spark 的arbitrary stateful streaming 来实现您自己的聚合;以及使用[flat]MapWithGroupState 在单个流上根据需要进行尽可能多的聚合。这将是一种合法且受支持的方式,与上述仅在某些情况下有效的黑客方式不同。此方法仅适用于 Spark Scala API,不适用于 PySpark。

  1. (也许有一天这将是一个长期的解决方案)

正确的方法是在 Spark Streaming 中显示对本机多重聚合的一些支持 - https://github.com/apache/spark/pull/23576 - 对此 SPARK jira/ PR 投票,如果您对此感兴趣,请表示您的支持。

【讨论】:

    【解决方案2】:

    从 spark 结构化流 2.4.5 开始,无状态处理不支持多个聚合。 但是如果您需要有状态的处理,可以多次聚合。

    通过追加模式,您可以在分组数据集(通过使用groupByKey API 获得)上多次使用flatMapGroupWithState API。

    【讨论】:

      【解决方案3】:

      在 Spark 2.4.4(目前最新)中不支持多流聚合,您可以使用 .foreachBatch() method

      一个虚拟的例子:

      query =  spark
              .readStream
              .format('kafka')
              .option(..)
              .load()
      
             .writeStream
             .trigger(processingTime='x seconds')
             .outputMode('append')
             .foreachBatch(foreach_batch_function)
             .start()
      
      query.awaitTermination()        
      
      
      def foreach_batch_function(df, epoch_id):
           # Transformations (many aggregations)
           pass   
      

      【讨论】:

        【解决方案4】:

        从 Spark 2.4 开始,不支持 Spark 结构化流中的多个聚合。支持这一点可能很棘手,尤其是。事件时间处于“更新”模式,因为聚合输出可能会随着迟到的事件而改变。在“附加”模式下支持这一点非常简单,但是 spark 还不支持真正的水印。

        建议以“追加”模式添加它 - https://github.com/apache/spark/pull/23576

        如果有兴趣,您可以观看 PR 并在那里投出您的投票。

        【讨论】:

          【解决方案5】:

          您没有提供任何代码,所以我将使用引用 here 的示例代码。

          假设下面是我们供 DF 使用的初始代码。

          import pyspark.sql.functions as F
          spark = SparkSession. ...
          
          # Read text from socket
          socketDF = spark \
              .readStream \
              .format("socket") \
              .option("host", "localhost") \
              .option("port", 9999) \
              .load()
          
          socketDF.isStreaming()    # Returns True for DataFrames that have streaming sources
          
          socketDF.printSchema()
          
          # Read all the csv files written atomically in a directory
          userSchema = StructType().add("name", "string").add("age", "integer")
          csvDF = spark \
              .readStream \
              .option("sep", ";") \
              .schema(userSchema) \
              .csv("/path/to/directory")  # Equivalent to format("csv").load("/path/to/directory")
          

          此处按 name 对 df 进行分组,并应用聚合函数 count、sum 和 balance。

          grouped = csvDF.groupBy("name").agg(F.count("name"), F.sum("age"), F.avg("age"))
          

          【讨论】:

          • 你能解释一下如何将数据帧“分组”写入蜂巢表吗?
          • 如果我做df_joined.writeStream.queryName('joined_query').outputMode('complete').format('memory').start(),我会在第一个问题中得到同样的错误:Multiple streaming aggregations are not supported with streaming DataFrame
          【解决方案6】:

          对于 spark 2.2 及更高版本(不确定早期版本),如果您可以将聚合设计为使用 flatMapGroupWithState 和 append 模式,则可以进行尽可能多的聚合你要。 这里提到了限制Spark structured streaming - Output mode

          【讨论】:

            【解决方案7】:

            这是不支持的,但也有其他方法。就像执行单个聚合并将其保存到 kafka 一样。从 kafka 中读取并再次应用聚合。这对我有用。

            【讨论】:

            • 我一直在尝试以这种方式解决类似的问题,但在验证行为时遇到了一些问题。如果我有两个应用程序,一个生成第一个聚合流(以更新模式输出),第二个读取流并执行第二个聚合;第二个应用程序会不断从第一个应用程序获取更新的值吗?
            • 您可能需要考虑到在 Spark 中写入 Kafka 并不是一次性的。 (有状态的完全一次和端到端的完全一次是不同的。)您可能会在输出主题中出现重复项,这可能会使第二次聚合不正确。这种方法必须以“端到端”方式执行一次。
            • @JungtaekLim,感谢我们测试它,我们实现它时没有得到任何重复
            • 如果其中一个分区写入失败而其他分区成功写入,则会提供重复,最终该批处理失败。当然,在愉快的情况下不会发生重复。
            【解决方案8】:

            这在 Spark 2.0 中不受支持,因为结构化流 API 仍处于试验阶段。请参阅 here 以查看所有当前限制的列表。

            【讨论】:

            • 我正在检查这个。我想它会起作用的。谢谢!
            • 由于缺乏结构化流 API 的支持,目前看来这是要走的路。
            猜你喜欢
            • 1970-01-01
            • 2021-04-24
            • 1970-01-01
            • 1970-01-01
            • 2020-09-12
            • 1970-01-01
            • 2020-03-19
            • 1970-01-01
            • 1970-01-01
            相关资源
            最近更新 更多