How to control Spark Stream with stream of dynamic queries?

Viewed 50

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?

0 Answers
Related