Pushkar Paranjpe
01/16/2025, 12:53 PMColin Ho
01/16/2025, 6:51 PMinto_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.explainjay
01/16/2025, 8:03 PMcount_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?Pushkar Paranjpe
01/17/2025, 10:44 AMPushkar Paranjpe
02/08/2025, 8:05 AMinto_partitions or repartitionjay
02/09/2025, 9:50 PMPushkar Paranjpe
03/06/2025, 3:52 AMPushkar Paranjpe
03/06/2025, 3:53 AMjay
03/06/2025, 3:54 AMPushkar Paranjpe
03/06/2025, 3:55 AMjay
03/06/2025, 3:55 AMPushkar Paranjpe
03/06/2025, 3:56 AMPushkar Paranjpe
03/06/2025, 3:56 AMPushkar Paranjpe
03/06/2025, 4:01 AMPushkar Paranjpe
03/06/2025, 4:01 AMPushkar Paranjpe
03/06/2025, 4:05 AMjay
03/06/2025, 4:05 AMread_parquet -> select(*subset_of_columns) -> write_parquet?Pushkar Paranjpe
03/06/2025, 4:06 AMjay
03/06/2025, 4:08 AMjay
03/06/2025, 4:08 AMPushkar Paranjpe
03/06/2025, 4:14 AMPushkar Paranjpe
03/06/2025, 4:57 AMRepartition: 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 = 2235548783203jay
03/06/2025, 5:00 AMPushkar Paranjpe
03/06/2025, 5:02 AMPushkar Paranjpe
03/06/2025, 5:03 AMjay
03/06/2025, 5:04 AMPushkar Paranjpe
03/06/2025, 5:05 AMmin_scan_task ? daft.set_execution_config ?jay
03/06/2025, 5:27 AMPushkar Paranjpe
03/06/2025, 5:31 AM* TabularScan:
| Num Scan Tasks = 140
| Estimated Scan Bytes = 2235548783203
after setting this exec config:
daft.set_execution_config(
scan_tasks_min_size_bytes=1,
)jay
03/06/2025, 5:32 AMjay
03/06/2025, 5:34 AMdaft.set_execution_config(
scan_tasks_min_size_bytes=1,
scan_tasks_max_size_bytes=1,
)Pushkar Paranjpe
03/06/2025, 5:36 AM* 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 }jay
03/06/2025, 5:41 AMPushkar Paranjpe
03/06/2025, 5:41 AMjay
03/06/2025, 6:07 AMdaft.set_execution_config(
max_sources_per_scan_task=1,
)
This should completely disable the merging.jay
03/06/2025, 6:07 AMPushkar Paranjpe
03/06/2025, 6:11 AM* TabularScan:
| Num Scan Tasks = 1400
| Estimated Scan Bytes = 2235548783203
| Clustering spec = { Num partitions = 1400 }Pushkar Paranjpe
03/06/2025, 6:17 AM(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.Pushkar Paranjpe
03/06/2025, 6:37 AMPushkar Paranjpe
03/06/2025, 6:37 AM