I have a stream of data coming from on Kafka call it a SourceStream.
I have another stream of Spark SQL queries whose individual values are Spark SQL queries along with a window size.
I want those queries to be applied to the SourceStream data, and pass the results of queries to the sink.
Eg. Source Stream
Id type timestamp user amount
------- ------ ---------- ---------- --------
uuid1 A 342342 ME 10.0
uuid2 B 234231 YOU 120.10
uuid3 A 234234 SOMEBODY 23.12
uuid4 A 234233 WHO 243.1
uuid5 C 124555 IT 35.12
...
....
Query Stream
Id window query
------- ------ ------
uuid13 1 hour select 'uuid13' as u, max(amount) as output from df where type = 'A' group by ..
uuid21 5 minute select 'uuid121' as u, count(1) as output from df where amount > 100 group by ..
uuid321 1 day select 'uuid321' as u, sum(amount) as output from df where amount > 100 group by ..
...
....
Each query in query stream would be applied to the source stream's incoming data at window mentioned along with the query, and the output would be sent to the sink.
What ways can I implement it with the Spark?