Is there a good way to change spark streaming filter at runtime?

Viewed 61

I’m using Spark 3.0.2 (with Scala) and I'm trying to figure out how I can achieve dynamic filtering (by a certain field that can be changed dynamically) while streaming from a source.

what I’ve tried:

*I thought it should be easy because spark streaming is based on micro batches so I can update the filter value on a separate thread but it seems that spark doesn’t work that way. it only “running/planning” the dataframe query/filter part once before the stream starts so any change to a variable/broadcast has no effect.

*join with another external source (as dataframe) that holds the relevant values to filter - works, but then the values are not a part of the filter push down and much more data is collected and filtered in the memory instead of before (source is parquet/delta io).

*using UDF in the filter - because it's a black box there will be no filter pushdown.

*stopping and restarting the streaming - may work but I’m not sure that it’s the best thing to do.

any experience with that will be helpful

0 Answers
Related