Details
-
Improvement
-
Status: Resolved
-
Major
-
Resolution: Fixed
-
Impala 2.8.0
Description
At the moment, runtime filters are sent from each fragment instance directly to the coordinator for aggregation (ORing) at the coordinator.
With multi-threaded execution, we will have an order of magnitude more fragment instances per node, at which point the coordinator would become a bottleneck during the aggregation process. To avoid that, we need to aggregate the local instances' runtime filters at each node before sending the filter off to the coordinator.
Attachments
Issue Links
- breaks
-
IMPALA-9395 RuntimeFilter causes impalad crash
- Resolved
- is blocked by
-
IMPALA-4063 Make fragment instance reports per-query (or per-host) instead of per-fragment instance.
- Resolved