him*_*ian 12 scala apache-spark spark-structured-streaming
我使用Spark 2.2.0-rc1.
我有一个卡夫卡topic这我查询运行水印聚集,有1 minute水印,发出来console与append输出模式.
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)
这是预期的行为吗?
向Kafka推送更多数据应该会触发Spark输出内容.目前的行为完全是因为内部实施.
当您推送一些数据时,StreamingQuery将生成一个要运行的批处理.当该批次完成时,它将记住该批次中的最大事件时间.然后在下一批中,因为您正在使用append模式,StreamingQuery将使用最大事件时间和水印来从StateStore中逐出旧值并输出它.因此,您需要确保生成至少两个批次才能查看输出.
这是我最好的猜测:
附加模式仅在水印通过后(例如在这种情况下 1 分钟后)输出数据。您没有设置触发器(例如.trigger(Trigger.ProcessingTime("10 seconds")),因此默认情况下它会尽可能快地输出批次。所以在第一分钟,你的所有批次都应该是空的,一分钟后的第一批应该包含一些内容。
另一种可能性是您正在使用groupBy("time")而不是groupBy(window("time", "[window duration]")). 我相信水印是为了与时间窗口或 mapGroupsWithState 一起使用,所以我不知道在这种情况下交互是如何工作的。