【发布时间】: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