Details
-
Improvement
-
Status: Closed
-
Critical
-
Resolution: Fixed
-
0.9.8
-
None
-
None
Description
Currently, incremental queries use Spark's coGroup to merge the current state with the results of processing the new data in the stream. With this patch, the merge is done with a special outer join that doesn't shuffle the state again (it only shuffles the results from the new data).
Attachments
Issue Links
- links to