Inconsistent spark structured streaming behavior when using ProcessingTimeTimeout

Viewed 100

Hope you can help.

The problem

Stack: spark 3.2.1, kafka.

After the first start of structured streaming query, which uses (flat)mapGroupsWithState function with GroupStateTimeout.ProcessingTimeTimeout(), empty batches are generated every trigger interval (e.g. 10 seconds) in case when there is no input data. This is ok and great actually, because processing time timeouts are processed.

However, after restart of the same query from checkpoint, empty batches are not generated every trigger interval in case when there is no input data. After data comes, empty batches are generated again. This is not ok. Here is a business case.

Imagine a processing time timeout was set to 12:00. Imagine app got down at 11:59 and was restarted at 12:01 (e.g. for maintenance). If new data hadn't arrived during this period, timeout would not be processed. It means app would contain expired state until new data arrives.

Reproducing the problem

  1. Start query (bootstrapServer, topic and checkpointLocation are parameters)
    val spark = SparkSession.builder().master("local[*]").getOrCreate()

    import spark.implicits._

    val query = spark.readStream
      .format("kafka")
      .option("kafka.bootstrap.servers", bootstrapServer)
      .option("subscribe", topic)
      .load()
      .select("value")
      .as[String]
      .groupByKey(v => v)
      .mapGroupsWithState(GroupStateTimeout.ProcessingTimeTimeout()) {
        (k: String, values: Iterator[String], state: GroupState[String]) =>
          state.setTimeoutDuration("1 minute")
          (k, values.size)
      }
      .writeStream
      .outputMode("update")
      .trigger(Trigger.ProcessingTime("10 seconds"))
      .option("checkpointLocation", checkpointLocation)
      .format("console")
      .option("truncate", "false")
      .start()

    query.awaitTermination()
  1. Ensure that empty batches are generated every trigger interval
-------------------------------------------
Batch: 0
-------------------------------------------
+---+---+
|_1 |_2 |
+---+---+
+---+---+

-------------------------------------------
Batch: 1
-------------------------------------------
+---+---+
|_1 |_2 |
+---+---+
+---+---+
  1. Write some record to topic
  2. Ensure it is processed
-------------------------------------------
Batch: 2
-------------------------------------------
+----+---+
|_1  |_2 |
+----+---+
|qwer|1  |
+----+---+
  1. Stop app
  2. Wait for more than 1 minute, i.e. timeout duration
  3. Start app again
  4. Wait for more than 10 seconds, i.e. trigger interval
  5. Ensure that empty batches are not generated
  6. Write another record to topic
  7. Ensure that timeout is processed
-------------------------------------------
Batch: 3
-------------------------------------------
+----+---+
|_1  |_2 |
+----+---+
|qwer|0  |
|asdf|1  |
+----+---+
  1. Ensure that empty batches are generated again
-------------------------------------------
Batch: 4
-------------------------------------------
+---+---+
|_1 |_2 |
+---+---+
+---+---+

-------------------------------------------
Batch: 5
-------------------------------------------
+---+---+
|_1 |_2 |
+---+---+
+---+---+

Research

Documentation

Documentation of GroupState - https://spark.apache.org/docs/3.2.1/api/scala/org/apache/spark/sql/streaming/GroupState.html:

With ProcessingTimeTimeout, the timeout duration can be set by calling GroupState.setTimeoutDuration. The timeout will occur when the clock has advanced by the set duration. Guarantees provided by this timeout with a duration of D ms are as follows:

  • Timeout will never occur before the clock time has advanced by D ms
  • Timeout will occur eventually when there is a trigger in the query (i.e. after D ms). So there is no strict upper bound on when the timeout would occur. For example, the trigger interval of the query will affect when the timeout actually occurs. If there is no data in the stream (for any group) for a while, then there will not be any trigger and timeout function call will not occur until there is data.
  • Since the processing time timeout is based on the clock time, it is affected by the variations in the system clock (i.e. time zone changes, clock skew, etc.).

This part seems incorrect:

If there is no data in the stream (for any group) for a while, then there will not be any trigger and timeout function call will not occur until there is data.

But as we've seen earlier, actually empty batches are generated (not always though) in case of data absence. As I checked the source code, documentation was written before related code was refactored.

Source code

I've also tried to find answers in spark source code, but all I've found is that code behaves as we've observed earlier. If you are curious, interesting part starts at this line.

Solution (crutch)

For now the solution is to write some garbage to the topic on the app start before starting structured streaming query.

Desired behavior

Query works the same way on the first start and after restart, ideally, always generating empty batches when there is no input data.

Finally, the question

Do you know any other ways (less crutchy, maybe) to trigger such query after restart?
If not, maybe you could tell the reason behind described behavior?
And also, would you agree with me that this behavior seems incorrect and desired behavior looks more consistent?

Thank you!

0 Answers
Related