Как подсчитать элементы за временное окно?

Я пытаюсь использовать структурированную потоковую передачу Spark для подсчета количества элементов из Kafka для каждого временного окна с помощью кода ниже:

import java.text.SimpleDateFormat
import java.util.Date
import org.apache.spark.sql.ForeachWriter
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions.window

object Counter extends App {
  val dateFormatter = new SimpleDateFormat("HH:mm:ss")
  val spark = ...
  import spark.implicits._

  val df = spark.readStream
    .format("kafka")
    .option("kafka.bootstrap.servers", ...)
    .option("subscribe", ...)
    .load()

  val windowDuration = "5 minutes"
  val counts = df
    .select("value").as[Array[Byte]]
    .map(decodeTimestampFromKafka).toDF("timestamp")
    .select($"timestamp" cast "timestamp")
    .withWatermark("timestamp", windowDuration)
    .groupBy(window($"timestamp", windowDuration, "1 minute"))
    .count()
    .as[((Long, Long), Long)]

  val writer = new ForeachWriter[((Long, Long), Long)] {
    var partitionId: Long = _
    var version: Long = _

    def open(partitionId: Long, version: Long): Boolean = {
      this.partitionId = partitionId
      this.version = version
      true
    }

    def process(record: ((Long, Long), Long)): Unit = {
      val ((start, end), docs) = record
      val startDate = dateFormatter.format(new Date(start))
      val endDate = dateFormatter.format(new Date(end))
      val now = dateFormatter.format(new Date)
      println(s"$now:$this|$partitionId|$version: ($startDate, $endDate) $docs")
    }

    def close(errorOrNull: Throwable): Unit = {}
  }

  val query = counts
    .repartition(1)
    .writeStream
    .outputMode("complete")
    .foreach(writer)
    .start()

  query.awaitTermination()

  def decodeTimestampFromKafka(bytes: Array[Byte]): Long = ...
}

Я ожидал, что раз в минуту (продолжительность слайда) он будет выводить одну запись (поскольку единственный ключ агрегации - это окно) с количеством элементов за последние 5 минут (продолжительность окна). Однако он выводит несколько записей 2-3 раза в минуту, как в этом примере:

...
22:44:34|Counter$$anon$1@6eb68dd7|0|8: (22:43:20, 22:43:20) 383
22:44:34|Counter$$anon$1@6eb68dd7|0|8: (22:43:18, 22:43:19) 435
22:44:34|Counter$$anon$1@6eb68dd7|0|8: (22:42:33, 22:42:34) 395
22:44:34|Counter$$anon$1@6eb68dd7|0|8: (22:43:14, 22:43:14) 435
22:44:34|Counter$$anon$1@6eb68dd7|0|8: (22:43:09, 22:43:09) 437
22:44:34|Counter$$anon$1@6eb68dd7|0|8: (22:43:19, 22:43:19) 411
22:44:34|Counter$$anon$1@6eb68dd7|0|8: (22:43:07, 22:43:07) 400
22:44:34|Counter$$anon$1@6eb68dd7|0|8: (22:43:17, 22:43:17) 392
22:44:44|Counter$$anon$1@5b70120f|0|9: (22:43:37, 22:43:38) 420
22:44:44|Counter$$anon$1@5b70120f|0|9: (22:43:25, 22:43:25) 395
22:44:44|Counter$$anon$1@5b70120f|0|9: (22:43:22, 22:43:22) 416
22:44:44|Counter$$anon$1@5b70120f|0|9: (22:43:00, 22:43:00) 438
22:44:44|Counter$$anon$1@5b70120f|0|9: (22:43:41, 22:43:41) 426
22:44:44|Counter$$anon$1@5b70120f|0|9: (22:44:13, 22:44:13) 132
22:44:44|Counter$$anon$1@5b70120f|0|9: (22:44:02, 22:44:02) 128
22:44:44|Counter$$anon$1@5b70120f|0|9: (22:44:09, 22:44:09) 120
...

Изменение режима вывода на append, похоже, меняет поведение, но все еще далеко от того, что я ожидал.

Что не так с моими предположениями о том, как это должно работать? Учитывая приведенный выше код, как следует интерпретировать или использовать образец вывода?


person erdavila    schedule 11.09.2017    source источник


Ответы (1)


Вы позволяете подсчитывать поздние события продолжительностью до 5 минут и обновлять окна, уже рассчитанные (withWatermark), см. обработка поздних данных и добавление водяных знаков в руководстве по Spark.

person Arnon Rotem-Gal-Oz    schedule 11.09.2017