【发布时间】:2018-10-28 01:31:20
【问题描述】:
我们正在使用结构化流对实时数据进行聚合。我正在创建一个可配置的 Spark 作业,该作业被赋予了配置,并使用它来对翻滚窗口中的行进行分组并执行聚合。我知道如何使用功能界面来做到这一点。
这是一个使用函数式接口的代码片段
var valStream = sparkSession.sql(sparkSession.sql(config.aggSelect)) //<- 1
.withWatermark("eventTime", "15 minutes") //<- 2
.groupBy(window($"eventTime", "1 minute"), $"aggCol1", $"aggCol2") //<- 3
.agg(count($"aggCol2").as("myAgg2Count"))
第 1 行执行来自配置的 SQL 字符串。我想将第 2 行和第 3 行移到 SQL 语法中,以便在配置中指定分组和聚合。
有没有人知道如何在 Spark SQL 中指定这个?
【问题讨论】:
标签: scala apache-spark spark-structured-streaming