附加模式下的水印聚合查询的空输出

him*_*ian 12 scala apache-spark spark-structured-streaming

我使用Spark 2.2.0-rc1.

我有一个卡夫卡topic这我查询运行水印聚集,有1 minute水印,发出来consoleappend输出模式.

import org.apache.spark.sql.types._
val schema = StructType(StructField("time", TimestampType) :: Nil)
val q = spark.
  readStream.
  format("kafka").
  option("kafka.bootstrap.servers", "localhost:9092").
  option("startingOffsets", "earliest").
  option("subscribe", "topic").
  load.
  select(from_json(col("value").cast("string"), schema).as("value"))
  select("value.*").
  withWatermark("time", "1 minute").
  groupBy("time").
  count.
  writeStream.
  outputMode("append").
  format("console").
  start
Run Code Online (Sandbox Code Playgroud)

我在Kafka推送以下数据topic:

{"time":"2017-06-07 10:01:00.000"}
{"time":"2017-06-07 10:02:00.000"}
{"time":"2017-06-07 10:03:00.000"}
{"time":"2017-06-07 10:04:00.000"}
{"time":"2017-06-07 10:05:00.000"}
Run Code Online (Sandbox Code Playgroud)

我得到以下输出:

scala> -------------------------------------------
Batch: 0
-------------------------------------------
+----+-----+                                                                    
|time|count|
+----+-----+
+----+-----+

-------------------------------------------
Batch: 1
-------------------------------------------
+----+-----+                                                                    
|time|count|
+----+-----+
+----+-----+

-------------------------------------------
Batch: 2
-------------------------------------------
+----+-----+                                                                    
|time|count|
+----+-----+
+----+-----+

-------------------------------------------
Batch: 3
-------------------------------------------
+----+-----+                                                                    
|time|count|
+----+-----+
+----+-----+

-------------------------------------------
Batch: 4
-------------------------------------------
+----+-----+                                                                    
|time|count|
+----+-----+
+----+-----+
Run Code Online (Sandbox Code Playgroud)

这是预期的行为吗?

zsx*_*ing 8

向Kafka推送更多数据应该会触发Spark输出内容.目前的行为完全是因为内部实施.

当您推送一些数据时,StreamingQuery将生成一个要运行的批处理.当该批次完成时,它将记住该批次中的最大事件时间.然后在下一批中,因为您正在使用append模式,StreamingQuery将使用最大事件时间和水印来从StateStore中逐出旧值并输出它.因此,您需要确保生成至少两个批次才能查看输出.


Ray*_*y J 5

这是我最好的猜测:

附加模式仅在水印通过后(例如在这种情况下 1 分钟后)输出数据。您没有设置触发器(例如.trigger(Trigger.ProcessingTime("10 seconds")),因此默认情况下它会尽可能快地输出批次。所以在第一分钟,你的所有批次都应该是空的,一分钟后的第一批应该包含一些内容。

另一种可能性是您正在使用groupBy("time")而不是groupBy(window("time", "[window duration]")). 我相信水印是为了与时间窗口或 mapGroupsWithState 一起使用,所以我不知道在这种情况下交互是如何工作的。

  • 我在尝试按时间变量分组时遇到了类似的问题。但是将其更改为 window(time, ...) 解决了我的问题 (2认同)