【问题标题】:How can I use aggregate with join in the same query result with Spark?如何在与 Spark 相同的查询结果中使用聚合和连接?
【发布时间】:2021-06-27 06:33:14
【问题描述】:

我需要加入以使用 postgres 数据丰富我的数据框。在 Spark Streaming 中我可以正常进行,因为数据是批量处理的。但是,在结构化流式传输中,每当我尝试将聚合与连接一起使用时都会出现错误。

例如:如果我在输出模式完成的情况下使用聚合,则作业可以正常工作,但是,如果我添加连接,则会返回错误:

Join between two streaming DataFrames/Datasets is not supported in Complete output mode, only in Append output mode;

如果我反其道而行之,也会发生同样的情况。当我使用输出模式追加连接时,作业正常运行,但是,如果我添加聚合,作业返回错误:

Append output mode not supported when there are streaming aggregations on streaming DataFrames/DataSets without watermark;

最后,我想知道是否有任何方法可以一起使用连接和聚合,而不会在使用 spark 结构化流时出错。

如果是这样,不产生这种错误的实现是什么样的?

def main(args: Array[String]): Unit = {
  val ss: SparkSession = Spark.getSparkSession
  val postgresSQL = new PostgresConnection
  val dataCollector = new DataCollector(postgresSQL)
  val collector = new Collector(ss,dataCollector)
  
  import ss.implicits._
  
  val stream: DataFrame = Kafka.setStructuredStream(ss)

  val parsed: DataFrame = Stream.parseInputMessages(stream)

  val getRelation: DataFrame = collector.getLastRelation(parsed)

    getRelation
    .writeStream.format("console")
    .trigger(Trigger.ProcessingTime(5000))
      .outputMode("complete")
      .queryName("Join")
      .start()

  ss.streams.awaitAnyTermination()
}

在我的 getLastRelation 方法中,我调用了 convertData 方法和 compareData 方法。

    def getLastRelation(messageToProcess: DataFrame): DataFrame = {
    // Faz tratamentos no DF para preparar a busca
    val dss: Dataset[Message] = this.convertData(messageToProcess)
    val dsRelacaolista: Dataset[WithStructure] = this.getPersonStructure(ds)
    val compareData = this.compareData(dsRelacaolista,messageToProcess)
    compareData
}

在我的 convertData 方法中,我使用了一个 agg。

    def convertData(data: DataFrame): Dataset[Message] = {
    data.selectExpr("country","code","order")
      .groupBy($"country",$"order")
      .agg(collect_list("code")
        .as("code"))
      .as[Message]
  }

在我的 compareData 方法中,我使用了 join:

def compareData(data: Dataset[WithStructure], message: DataFrame): DataFrame = {
    val tableJoin = message.selectExpr("order","order_id","hashCompare","created_at")

    data.toDF()
      .withColumn("hashCompare",hash($"country",$"code"))
      .join(tableJoin,"hashCompare")
  }

注意:我使用 scala 作为一种语言(我不知道这些信息对这个问题是否重要)

【问题讨论】:

  • 这就是它的实现方式。有限制

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


【解决方案1】:

如果您对流进行聚合查询,则需要指定水印和窗口。

例如:

 data
     .withWatermark("created_at", "10 minutes")
     .selectExpr("country","code","order")
     .groupBy(window($"created_at", "10 minutes", "5 minutes"), $"country",$"order")
     .agg(collect_list("code")
     .as("code"))
     .as[Message]

随流到达的数据可能因任何原因延迟(由于网络速度变慢等)。水印允许指定聚合应等待滞后事件多长时间。所有延迟高于水印中指定时间的事件都将被忽略。

追加模式不允许修改之前输出的结果。因此,它需要水印来确保聚合数据不会被进一步更新。

您可以为水印选择更长的窗口,这将使您对处理延迟数据的容忍度更高。缺点是上游会被 watermark 持续时间延迟,因为查询必须等待 watermark 中指定的时间过去,然后才能完成聚合。

此外,对于流-流连接(当连接的两边都是流数据集时),您还需要指定窗口。

来自文档:

在两个数据流之间生成连接结果的挑战在于,在任何时间点,连接两侧的数据集视图都不完整,因此更难找到输入之间的匹配。

请查看有关 stream-stream joinsdoing streaming aggregations 的文档。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2011-02-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多