【发布时间】:2018-08-19 20:47:27
【问题描述】:
我有以下代码,它输出
number: 1, count: 1
number: 2, count: 1
number: 3, count: 2
number: 6, count: 2
number: 7, count: 1
我认为number: 6, count: 2 不应该输出,因为事件低于水位线。但我不明白它为什么输出
import java.sql.Timestamp
import org.apache.spark.sql.execution.streaming.MemoryStream
import org.apache.spark.sql.{ForeachWriter, Row, SparkSession}
object UpdateModeWithWatermarkTest {
def main(args: Array[String]): Unit = {
val spark: SparkSession = SparkSession.builder()
.appName("UpdateModeWithWatermarkTest")
.config("spark.sql.shuffle.partitions", 1)
.master("local[2]").getOrCreate()
import spark.implicits._
val inputStream = new MemoryStream[(Timestamp, Int)](1, spark.sqlContext)
val now = 5000L
val aggregatedStream = inputStream.toDS().toDF("created", "number")
.withWatermark("created", "1 second")
.groupBy("number")
.count()
val query = aggregatedStream.writeStream.outputMode("update")
.foreach(new ForeachWriter[Row] {
override def open(partitionId: Long, epochId: Long): Boolean = true
override def process(value: Row): Unit = {
println(s"number: ${value.getInt(0)}, count: ${value.getLong(1)}")
}
override def close(errorOrNull: Throwable): Unit = {}
}).start()
new Thread(new Runnable() {
override def run(): Unit = {
inputStream.addData(
(new Timestamp(now + 5000), 1),
(new Timestamp(now + 5000), 2),
(new Timestamp(now + 5000), 3),
(new Timestamp(now + 5000), 3)
)
while (!query.isActive) {
Thread.sleep(50)
}
Thread.sleep(10000)
// At this point, the water mark is (now + 5000) - 1 second = 9 seconds
// when adding following two events: (new Timestamp(4000L), 6), (new Timestamp(now), 6)
// These two events are below water mark, so that they should be discarded, then should not output number: 6, count: 2
inputStream.addData((new Timestamp(4000L), 6))
inputStream.addData(
(new Timestamp(now), 6),
(new Timestamp(11000), 7)
)
}
}).start()
query.awaitTermination(45000)
}
}
【问题讨论】:
-
你在这个 Tom 上的表现如何?
标签: apache-spark spark-structured-streaming