【问题标题】:How to specify structured streaming time based window in straight Spark SQL如何在直接 Spark SQL 中指定基于结构化流时间的窗口
【发布时间】: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


    【解决方案1】:

    withWatermark 没有对应的 SQL 语法。您必须使用数据框 API。

    对于聚合,您可以执行类似的操作

    select count(aggcol2) as myAgg2Count
    from xxx
    group by window(eventTime, "1 minute"), aggCo1, aggCol2
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2019-10-29
      • 2017-10-09
      • 2018-07-20
      • 1970-01-01
      • 1970-01-01
      • 2019-10-10
      • 1970-01-01
      相关资源
      最近更新 更多