Suppose I want to run a streaming job that takes new data every x seconds and outputs new rows for each trigger without any aggregation. For example:
val query = wordCounts.writeStream
.outputMode("append")
.trigger(Trigger.ProcessingTime("15 seconds"))
.format("console")
.start()
In the Spark documentation, I see the following:
"A query on the input will generate the “Result Table”. Every trigger interval (say, every 1 second), new rows get appended to the Input Table, which eventually updates the Result Table. Whenever the result table gets updated, we would want to write the changed result rows to an external sink."
"Append Mode - Only the new rows appended in the Result Table since the last trigger will be written to the external storage. This is applicable only on the queries where existing rows in the Result Table are not expected to change."
I understand that in Append mode, only new rows are written to the console, however the "Result Table" will still have historical rows from previous triggers correct? So, this table size could be increasing in size for a long time. Is there a way to delete these old rows from the result table since they are not needed and will not be changed?
I see the windowing + watermarking feature, however that is used mainly for aggregation queries right?