Hi daft team! First- thanks for this cool product....
# general
p
Hi daft team! First- thanks for this cool product. I tried this out on a dataset of ~100B records. Found good success out of the box with groupby aggs such as count and sum. But ran into issues with groupby count_distinct / approx_count_distinct. Ray dashboard shows memory getting swamped immediately and major spillage to disk, workers dying due to memory pressure etc. • is this a know failure mode of daft ? • any remedies to try ? • Known benchmarks specifically with groupby count_distinct ?
c
Hey @Pushkar Paranjpe, thanks for trying out Daft! Good to know that count and sum worked for you. count_distinct is going to be memory intensive because it will coalesce all the data to a single node to perform the count_distinct operation. cc @Raunak Bhagat if you have any tips here. For approx_count_distinct, you can try increase the number of partitions via
into_partitions
or
repartition
before the groupby. Higher number of partitions will also mean smaller partitions, and reduce the chance of OOMs. Also, do you know which stage the crashes happen? You should be able to see them under the jobs view in the Ray Dashboard. Lastly, if you're able to, it'd be great if you could show us the query plan,
df.explain(True)
! https://www.getdaft.io/projects/docs/en/stable/api_docs/doc_gen/dataframe_methods/daft.DataFrame.explain.html#daft.DataFrame.explain
j
The exact
count_distinct
will indeed be very memory-intensive. I don't think there's a workaround there
approx_count_distinct
though should be relatively lightweight -- I'm surprised that it causes OOMs. Do you expect the cardinality of your groups (total number of groups after the groupby) to be extremely high?
p
Thanks! I will try above suggestions and revert back. High cardinality - ~200M
No success even after - increasing the number of partitions via
into_partitions
or
repartition
j
The query plan would be really useful here for us to get an idea of what your workload looks like!
p
Memory pressure triggers killing of nodes and job fails.
My stack is daft|ray|k8s
j
Is kubernetes is killing the nodes?
p
Yep
j
Could you share the query plan? Would be helpful to help debug
p
Tried a limit ramp, works fine upto a threshold number of rows
Lemme see how i can share that, considering it will contain bunch of proprietary info
🙌 1
Have tried everything from parquet inflation factor, morsel size, custom ioconfig limiting concurrent reads (default was overwhelming s3), intoPartitions, repartition, toRayDataset (with num_cpu over provisioning), etc
Input partition is large - 2.2TB having 2B rows
Also - this is a purely narrow transformation- consume input dataframe and then output a dataframe having a subset of columns; no aggregations nor any joins
j
Is it just
read_parquet -> select(*subset_of_columns) -> write_parquet
?
p
Pretty much! One exception being a col -> f(col) Where f is O(1)
j
The plan will also have some helpful information about the number of ScanTasks being produced -- we do to coalesce files if we think that the Parquet files are small enough to place into the same partition
Likely what's happening here is that we're underestimating the amount of memory it takes to read the data. That's my guess, but hard to tell without the plan...
👍 1
p
Some way to override the memory estimate?
Copy code
Repartition: Scheme = Random
|   Num partitions = Some(100)
|   Stats = { Approx num rows = 10,000,000, Approx size bytes = 3.07 GiB,
|     Accumulated selectivity = 0.03 }
|
* Limit: 10000000
|   Stats = { Approx num rows = 10,000,000, Approx size bytes = 3.07 GiB,
|     Accumulated selectivity = 0.03 }
|
* Num Scan Tasks = 1400
...
* TabularScan:
|   Num Scan Tasks = 201
|   Estimated Scan Bytes = 2235548783203
j
Ah is the number of scan tasks going from 1400 to 201 after optimization?
p
yeah
the source partition contains 1400 part parquet files
j
Ok yeah pretty sure that’s the problem. It’s getting coalesced too aggressively. You can try turning off the coalescing via the execution config for min_scan_task size. Just make the minimum like 1 or something Then when you run explain again, check that the number of final scan tasks is 1400?
p
thanks, where do i set
min_scan_task
? daft.set_execution_config ?
p
Getting this in the plan:
Copy code
* TabularScan:
|   Num Scan Tasks = 140
|   Estimated Scan Bytes = 2235548783203
after setting this exec config:
Copy code
daft.set_execution_config(
    scan_tasks_min_size_bytes=1,
)
j
To confirm -- it's now going from 1400 scan tasks to 140? That's very odd
Can you try setting
Copy code
daft.set_execution_config(
    scan_tasks_min_size_bytes=1,
    scan_tasks_max_size_bytes=1,
)
👍 1
p
Copy code
* Num Scan Tasks = 1400
...
* Limit: 10000000
|   Eager = false
|   Num partitions = 140
|
* TabularScan:
|   Num Scan Tasks = 140
|   Estimated Scan Bytes = 2235548783203
|   Clustering spec = { Num partitions = 140 }
j
Ok it looks like something is really off wrt how the ScanTasks are being merged here I'm guessing. Let me bring this to the attention of the team tomorrow and see if we can provide a solution
p
thank you 👍
j
Ah can you try this?
Copy code
daft.set_execution_config(
    max_sources_per_scan_task=1,
)
This should completely disable the merging.
Once we verify that the number of files is 1400 we can try to run the workload
👍 1
p
that ^ seems to have done it, plan:
Copy code
* TabularScan:
|   Num Scan Tasks = 1400
|   Estimated Scan Bytes = 2235548783203
|   Clustering spec = { Num partitions = 1400 }
🙌 1
facing the same issue upon running the plan though :
Copy code
(raylet) The node with node id: ... and address: ... and node name: ... has been marked dead because the detector has missed too many heartbeats from it. This can happen when a 	(1) raylet crashes unexpectedly (OOM, etc.) 
	(2) raylet has lagging heartbeats due to slow network or busy workload.
Node memory fills up and it is killed. Multiple nodes die and then the job fails.
Executes successfully on a .limit(1M) dataset in ~2mins
Fails in above failure mode for .limit(>=10M) dataset size