Perhaps the most elegant way to express this is with the cumulative_eval expression:
(
df
.with_column(
pl.col("c1")
.cumulative_eval(pl.element().list())
.alias("c4")
)
)
shape: (4, 4)
┌─────┬─────┬─────┬───────────────┐
│ c1 ┆ c2 ┆ c3 ┆ c4 │
│ --- ┆ --- ┆ --- ┆ --- │
│ i64 ┆ i64 ┆ i64 ┆ list[i64] │
╞═════╪═════╪═════╪═══════════════╡
│ 1 ┆ 2 ┆ 3 ┆ [1] │
├╌╌╌╌╌┼╌╌╌╌╌┼╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 3 ┆ 4 ┆ 5 ┆ [1, 3] │
├╌╌╌╌╌┼╌╌╌╌╌┼╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 5 ┆ 6 ┆ 7 ┆ [1, 3, 5] │
├╌╌╌╌╌┼╌╌╌╌╌┼╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 7 ┆ 8 ┆ 9 ┆ [1, 3, ... 7] │
└─────┴─────┴─────┴───────────────┘
Some Notes
In general, I would not recommend that the cumulative_eval be used in this manner. As the documentation says:
Warning: This can be really slow as it can have O(n^2) complexity. Don’t use this for operations that visit all elements.
That said, in this case you are more likely to run out of system RAM long before this becomes a serious performance bottleneck.
For example, if we expand the size of the dataset to a mere 100,000 rows:
stop = 200_000
df = pl.DataFrame(
{
"c1": pl.arange(1, stop, step=2, eager=True),
"c2": pl.arange(2, stop + 1, step=2, eager=True),
"c3": pl.arange(3, stop + 2, step=2, eager=True),
}
)
df
shape: (100000, 3)
┌────────┬────────┬────────┐
│ c1 ┆ c2 ┆ c3 │
│ --- ┆ --- ┆ --- │
│ i64 ┆ i64 ┆ i64 │
╞════════╪════════╪════════╡
│ 1 ┆ 2 ┆ 3 │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┤
│ 3 ┆ 4 ┆ 5 │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┤
│ 5 ┆ 6 ┆ 7 │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┤
│ 7 ┆ 8 ┆ 9 │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┤
│ ... ┆ ... ┆ ... │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┤
│ 199993 ┆ 199994 ┆ 199995 │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┤
│ 199995 ┆ 199996 ┆ 199997 │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┤
│ 199997 ┆ 199998 ┆ 199999 │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┤
│ 199999 ┆ 200000 ┆ 200001 │
└────────┴────────┴────────┘
And look at the performance in terms of time and system RAM required:
import time
start = time.perf_counter()
result = (
df
.with_column(
pl.col("c1")
.cumulative_eval(pl.element().list())
.alias("c4")
)
)
result
print(time.perf_counter() - start)
print(result.estimated_size() / (2**30))
shape: (100000, 4)
┌────────┬────────┬────────┬────────────────────┐
│ c1 ┆ c2 ┆ c3 ┆ c4 │
│ --- ┆ --- ┆ --- ┆ --- │
│ i64 ┆ i64 ┆ i64 ┆ list[i64] │
╞════════╪════════╪════════╪════════════════════╡
│ 1 ┆ 2 ┆ 3 ┆ [1] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 3 ┆ 4 ┆ 5 ┆ [1, 3] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 5 ┆ 6 ┆ 7 ┆ [1, 3, 5] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 7 ┆ 8 ┆ 9 ┆ [1, 3, ... 7] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ ... ┆ ... ┆ ... ┆ ... │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 199993 ┆ 199994 ┆ 199995 ┆ [1, 3, ... 199993] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 199995 ┆ 199996 ┆ 199997 ┆ [1, 3, ... 199995] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 199997 ┆ 199998 ┆ 199999 ┆ [1, 3, ... 199997] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 199999 ┆ 200000 ┆ 200001 ┆ [1, 3, ... 199999] │
└────────┴────────┴────────┴────────────────────┘
>>> print(time.perf_counter() - start)
27.60901775200182
>>> print(result.estimated_size() / (2**30))
37.25625575333834
The algorithm took about 28 seconds to complete ... but creates a dataset that is 37 GB in RAM. As such, you are more likely to run out of system RAM long before the wall-clock performance becomes an issue.
Also note that the cumuluative_eval expression is marked as experimental and may be subject to future change.
Edit: columns of type i32
After looking at my results, I realized that I created the columns of type i64, rather than i32 (as in your example).
Nonetheless, the general results hold. For example, using 100,000 records of type i32, the results are:
shape: (100000, 4)
┌────────┬────────┬────────┬────────────────────┐
│ c1 ┆ c2 ┆ c3 ┆ c4 │
│ --- ┆ --- ┆ --- ┆ --- │
│ i32 ┆ i32 ┆ i32 ┆ list[i32] │
╞════════╪════════╪════════╪════════════════════╡
│ 1 ┆ 2 ┆ 3 ┆ [1] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 3 ┆ 4 ┆ 5 ┆ [1, 3] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 5 ┆ 6 ┆ 7 ┆ [1, 3, 5] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 7 ┆ 8 ┆ 9 ┆ [1, 3, ... 7] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ ... ┆ ... ┆ ... ┆ ... │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 199993 ┆ 199994 ┆ 199995 ┆ [1, 3, ... 199993] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 199995 ┆ 199996 ┆ 199997 ┆ [1, 3, ... 199995] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 199997 ┆ 199998 ┆ 199999 ┆ [1, 3, ... 199997] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 199999 ┆ 200000 ┆ 200001 ┆ [1, 3, ... 199999] │
└────────┴────────┴────────┴────────────────────┘
>>> print(time.perf_counter() - start)
15.642635091004195
>>> print(result.estimated_size() / (2**30))
18.628500409424305
Using i32 lowers the wall-clock run-time to 16 seconds, and creates a dataset of about 19 GB in RAM.
Still, if we double the number of records to 200,000 records:
shape: (200000, 4)
┌────────┬────────┬────────┬────────────────────┐
│ c1 ┆ c2 ┆ c3 ┆ c4 │
│ --- ┆ --- ┆ --- ┆ --- │
│ i32 ┆ i32 ┆ i32 ┆ list[i32] │
╞════════╪════════╪════════╪════════════════════╡
│ 1 ┆ 2 ┆ 3 ┆ [1] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 3 ┆ 4 ┆ 5 ┆ [1, 3] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 5 ┆ 6 ┆ 7 ┆ [1, 3, 5] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 7 ┆ 8 ┆ 9 ┆ [1, 3, ... 7] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ ... ┆ ... ┆ ... ┆ ... │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 399993 ┆ 399994 ┆ 399995 ┆ [1, 3, ... 399993] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 399995 ┆ 399996 ┆ 399997 ┆ [1, 3, ... 399995] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 399997 ┆ 399998 ┆ 399999 ┆ [1, 3, ... 399997] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 399999 ┆ 400000 ┆ 400001 ┆ [1, 3, ... 399999] │
└────────┴────────┴────────┴────────────────────┘
>>> print(time.perf_counter() - start)
45.241694554002606
>>> print(result.estimated_size() / (2**30))
74.50990379601717
The result takes 45 seconds of wall-time time, and creates a dataset that occupies almost 75 GB of system RAM.
Edit: using with_row_count
I understand that df.with_row_count() shall generate a Polars dataframe with extra column of increasing row IDs. Is it possible to use this row ID as a parameter to call pl.col('c1').head().to_list() to generate the same C4 column?
The problem is that the head expression does not take an Expression as a parameter, only an integer. That is, we cannot use a reference to another column as the n parameter in the head expression.
For example:
(
df
.with_row_count()
.with_column(
pl.col("c1")
.head(pl.col('row_nr'))
.list()
.alias("c4")
)
)
Traceback (most recent call last):
File "<stdin>", line 6, in <module>
File "/home/corey/.virtualenvs/StackOverflow3.10/lib/python3.10/site-packages/polars/internals/expr/expr.py", line 3219, in head
return wrap_expr(self._pyexpr.head(n))
TypeError: argument 'n': 'Expr' object cannot be interpreted as an integer
It is possible to get around this limitation by using apply/map. However, using apply or map typically leads to single-threaded, poor performance.
For example, we can use map and create a function that uses head and to_list to perform the action as follow:
start = time.perf_counter()
def head_func(s: pl.Series) -> pl.Series:
values = [s.head(nbr_recs + 1).to_list()
for nbr_recs in range(0, s.len())]
return pl.Series(values=values)
(
df
.with_column(
pl.col("c1")
.map(head_func)
.alias("c4")
)
)
print(time.perf_counter() - start)
shape: (100000, 4)
┌────────┬────────┬────────┬────────────────────┐
│ c1 ┆ c2 ┆ c3 ┆ c4 │
│ --- ┆ --- ┆ --- ┆ --- │
│ i32 ┆ i32 ┆ i32 ┆ list[i64] │
╞════════╪════════╪════════╪════════════════════╡
│ 1 ┆ 2 ┆ 3 ┆ [1] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 3 ┆ 4 ┆ 5 ┆ [1, 3] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 5 ┆ 6 ┆ 7 ┆ [1, 3, 5] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 7 ┆ 8 ┆ 9 ┆ [1, 3, ... 7] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ ... ┆ ... ┆ ... ┆ ... │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 199993 ┆ 199994 ┆ 199995 ┆ [1, 3, ... 199993] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 199995 ┆ 199996 ┆ 199997 ┆ [1, 3, ... 199995] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 199997 ┆ 199998 ┆ 199999 ┆ [1, 3, ... 199997] │
├╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌┼╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌╌┤
│ 199999 ┆ 200000 ┆ 200001 ┆ [1, 3, ... 199999] │
└────────┴────────┴────────┴────────────────────┘
>>> print(time.perf_counter() - start)
243.5193853039964
Problem is, for 100,000 records, our run-time went from about 16 seconds to 244 seconds.