dask / dask/dask-expr

Parquet statistics are collected twice

Open
#365 10 comments 0 reactions 0 assignees View on GitHub
Dominant language
Python
Stars
89
Forks
26
PR merge metrics
No merged PRs in 30d

Description

Currently, the parquet reader is using `split_row_groups='infer'` to infer whether it should split a file into multiple dask tasks or not which is done by collecting parquet statistics from every file.

Running the code below triggers this statistics collection twice. Once during the definition of the computation / while I'm mutating the dataframe and once as soon as I call compute.

```python
import dask_expr as dd

VAR1 = datetime(1998, 9, 2)

lineitem_ds = dd.read_parquet("s3://coiled-runtime-ci/tpc-h/scale-1000/lineitem")

lineitem_filtered = lineitem_ds[lineitem_ds.l_shipdate <= VAR1]
lineitem_filtered["sum_qty"] = lineitem_filtered.l_quantity
lineitem_filtered["sum_base_price"] = lineitem_filtered.l_extendedprice
lineitem_filtered["avg_qty"] = lineitem_filtered.l_quantity
lineitem_filtered["avg_price"] = lineitem_filtered.l_extendedprice

# This line now triggers a statistics collection iff `split_row_groups` is not False
lineitem_filtered["sum_disc_price"] = lineitem_filtered.l_extendedprice * (
1 - lineitem_filtered.l_discount
)

lineitem_filtered["sum_charge"] = (
lineitem_filtered.l_extendedprice
* (1 - lineitem_filtered.l_discount)
* (1 + lineitem_filtered.l_tax)
)

lineitem_filtered["avg_disc"] = lineitem_filtered.l_discount
lineitem_filtered["count_order"] = lineitem_filtered.l_discount

gb = lineitem_filtered.groupby(["l_returnflag", "l_linestatus"])

total = gb.agg(
{
"sum_qty": "sum",
"sum_base_price": "sum",
"sum_disc_price": "sum",
"sum_charge": "sum",
"avg_qty": "mean",
"avg_price": "mean",
"avg_disc": "mean",
"count_order": "count",
}
)

# Once I compute, another stats collection is triggered
total.compute()
```

![image](https://github.com/dask-contrib/dask-expr/assets/8629629/ae45d60e-5a32-4bbf-8bfc-fe29c9c0c71e)

related to https://github.com/dask-contrib/dask-expr/issues/363

Contributor guide

Open the contributing guide

Research direction

Start by running the supplied dask_expr example and tracing the parquet reader's split_row_groups='infer' path during dataframe mutation and compute. Identify why statistics are collected in both phases; done when the same computation collects parquet statistics only once while preserving task-splitting behavior.

Written by the indexing model from the issue text.

Assessment

Tech stack
python
Domain
data-engineering
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.