I have a spark structured streaming, where I need to stop the streaming process once I have the correct records count.
So far I have a method to stop streaming at particular interval of no activity, how can I include a counter to get the count and close the stream?
val resultStream=furtherFlattening
.writeStream
.format("console")
.option("truncate","false")
.trigger(Trigger.ProcessingTime(5, TimeUnit.SECONDS))
//. trigger(Trigger.ProcessingTime(5, TimeUnit.MINUTES))
.start()
.awaitTermination()
def stopStreamQuery(query: StreamingQuery, awaitTerminationTimeMs: Long,spark:SparkSession) {
while (query.isActive) {
val msg = query.status.message
if (!query.status.isDataAvailable
&& !query.status.isTriggerActive
&& !msg.equals("Initializing sources")) {
query.stop()
spark.close()
}
query.awaitTermination(awaitTerminationTimeMs)
}
}