Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion crates/polars-python/src/lazyframe/visit.rs
Original file line number Diff line number Diff line change
Expand Up @@ -58,7 +58,7 @@ impl NodeTraverser {
// Increment major on breaking changes to the IR (e.g. renaming
// fields, reordering tuples), minor on backwards compatible
// changes (e.g. exposing a new expression node).
const VERSION: Version = (13, 0);
const VERSION: Version = (14, 0);

pub fn new(root: Node, lp_arena: Arena<IR>, expr_arena: Arena<AExpr>) -> Self {
Self {
Expand Down
4 changes: 2 additions & 2 deletions crates/polars-python/src/lazyframe/visitor/expr_nodes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1466,8 +1466,8 @@ pub(crate) fn into_py(py: Python<'_>, expr: &AExpr) -> PyResult<Py<PyAny>> {
IRFunctionExpr::RowDecode(..) => {
return Err(PyNotImplementedError::new_err("row_decode"));
},
IRFunctionExpr::DynamicPred { .. } => {
return Err(PyNotImplementedError::new_err("dynamic_pred"));
IRFunctionExpr::DynamicPred { pred } => {
("dynamic_pred", pred.id().map(|u| u.as_u128())).into_py_any(py)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

question: where does the predicate that one must evaluate actually live? It seems like it's in pred.pred but that is not exposed anywhere. I'm also not sure it can be because it's an Arc<dyn PredicateExpr>?

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The predicate is simple in this case. For example, col < threshold (top_k) which we know from the sort order.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is this something that should be exposed? I don't think the predicate can be called. In this case the predicate is simple, but that doesn't seem like a guarantee.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think we could reconstruct the filter (even for complex predicates) because the unique id lets us match each dynamic_pred to its parent sort-slice. But I would want to avoid this manual reconstruction if possible.

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I read through @ritchie46's PR but I think I didn't fully understand the structure of how these things are matched up.

IIUC, we have something like:

    df = pl.DataFrame({"x": [1], "y": [1]})
    plan = df.lazy().with_columns(pl.col.x * pl.col.x).sort("y").head(3)

That produces:

SORT BY [slice: (0, 3, dynamic_pred: id-1)] [col("y")]
   WITH_COLUMNS:
   [[(col("x")) * (col("x"))]] 
    FILTER col("y").dynamic_predicate() # id-1
    FROM
      DF ["x", "y"]; PROJECT */2 COLUMNS

And so the idea here is that you're going to read df in chunks and apply the filter based on things you've already seen. So you read the first chunk, the predicate initially return true until you've "filled up" your slice. Then the next time the predicate runs on the next chunk, it delivers values if they are less than the max in the filled up slice, and so forth.

OK, in that scenario I can see how we can have an expression for the predicate.

But what about (it's not implemented yet) if there was a transformation like:

df = pl.DataFrame({"x": [1], "y": [1]})
plan = df.lazy().with_columns(pl.col.x * pl.col.x).unique("y")

That produced:

UNIQUE BY [col("y"), dynamic_pred: id-1]
  WITH_COLUMNS:
    [[(col("x")) * (col("x"))]]
      FILTER col("y").dynamic_predicate(id-1)
      FROM
        ...

Where in this case the dynamic_predicate is set membership of the already seen values.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Indeed. That's how it might be used as well. Another thing we plan to use it for is members of the hash-table in a join. But this is something that will dynamically at runtime be determined. Not something we can statically make a predicate for. Otherwise we would have done that already. :)

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

OK so the id is sufficient today to reconstruct the predicate. But as you all said we could set membership or members for the hash table for a join. We could (TODO) tag the predicate with a description (eg. TopK, UniqueMembership, JoinMembership)?

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't think it should be tagged. The id should be sufficient. You would find a dynamicpredicate with id=x in a scan and then the dynamicpredicate setting with id=x in a different IR node. The IR node would indicate what to do. How that is done is left to the implementation.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Ok cool. I bumped the IR version in 8fcf646. Also gentle reminder about adding me as a code owner for the visitor in the PR description.

},
}?,
options: py.None(),
Expand Down
6 changes: 4 additions & 2 deletions crates/polars-python/src/lazyframe/visitor/nodes.rs
Original file line number Diff line number Diff line change
Expand Up @@ -228,7 +228,7 @@ pub struct Sort {
#[pyo3(get)]
sort_options: (bool, Vec<bool>, Vec<bool>),
#[pyo3(get)]
slice: Option<(i64, usize)>,
slice: Option<(i64, usize, Option<u128>)>,
}

#[pyclass(frozen)]
Expand Down Expand Up @@ -520,7 +520,9 @@ pub(crate) fn into_py(py: Python<'_>, plan: &IR) -> PyResult<Py<PyAny>> {
sort_options.nulls_last.clone(),
sort_options.descending.clone(),
),
slice: slice.as_ref().map(|t| (t.0, t.1)),
slice: slice
.as_ref()
.map(|t| (t.0, t.1, t.2.as_ref().map(|p| p.id().as_u128()))),
}
.into_py_any(py),
IR::Cache { input, id } => Cache {
Expand Down
Loading