-
Notifications
You must be signed in to change notification settings - Fork 143
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
[FEAT] [New Query Planner] Support for Ray runner in new query planne…
…r. (#1265) This PR enables the Ray runner in the new query planner. This is accomplished by encapsulating the Rust-side `PhysicalPlan` in a Python-facing `PhysicalPlanScheduler`, which we then make pickleable. ## TODOs - [x] Add support for `ResourceRequest`s
- Loading branch information
1 parent
5a2dc7a
commit c43c76c
Showing
51 changed files
with
440 additions
and
227 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,5 +1,5 @@ | ||
from __future__ import annotations | ||
|
||
from daft.planner.planner import QueryPlanner | ||
from daft.planner.planner import PhysicalPlanScheduler | ||
|
||
__all__ = ["QueryPlanner"] | ||
__all__ = ["PhysicalPlanScheduler"] |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,13 +1,13 @@ | ||
from __future__ import annotations | ||
|
||
from daft.daft import LogicalPlanBuilder as _LogicalPlanBuilder | ||
from daft.daft import PhysicalPlanScheduler as _PhysicalPlanScheduler | ||
from daft.execution import physical_plan | ||
from daft.planner.planner import PartitionT, QueryPlanner | ||
from daft.planner.planner import PartitionT, PhysicalPlanScheduler | ||
|
||
|
||
class RustQueryPlanner(QueryPlanner): | ||
def __init__(self, builder: _LogicalPlanBuilder): | ||
self._builder = builder | ||
class RustPhysicalPlanScheduler(PhysicalPlanScheduler): | ||
def __init__(self, scheduler: _PhysicalPlanScheduler): | ||
self._scheduler = scheduler | ||
|
||
def plan(self, psets: dict[str, list[PartitionT]]) -> physical_plan.MaterializedPhysicalPlan: | ||
return physical_plan.materialize(self._builder.to_partition_tasks(psets)) | ||
def to_partition_tasks(self, psets: dict[str, list[PartitionT]]) -> physical_plan.MaterializedPhysicalPlan: | ||
return physical_plan.materialize(self._scheduler.to_partition_tasks(psets)) |
Oops, something went wrong.