Neil Wadhvana
02/12/2025, 12:16 AMdf.into_partitions(df.count_rows())
a. Outcome: late-job OOM 1.5-2 hours into the job. Makes it rather difficult to recover since validating the binary data for corruption would require unpacking the original data from the npz files anyways
Are there any concurrency knobs I can tune to limit the number of parallel executed tasks? I think I can manually start the ray head node with less CPUs than are available on the instance, and still keep the number of partitions very high? Any other ideas?Neil Wadhvana
02/12/2025, 1:05 AMNeil Wadhvana
02/12/2025, 1:06 AM(raylet) [2025-02-12 00:59:06,030 E 386983 386983] <http://node_manager.cc:3178|node_manager.cc:3178>: 7 Workers (tasks / actors) killed due to memory pressure (OOM), 0 Workers crashed due to other reasons at node (ID: 325e5979bf4e2a0d893fea1af94b453674aa37c566d739a5f643858e, IP: {REDACTED}) over the last time period. To see more information about the Workers killed on this node, use `ray logs raylet.out -ip {REDACTED}`
(raylet) Refer to the documentation on how to address the out of memory issue: <https://docs.ray.io/en/latest/ray-core/scheduling/ray-oom-prevention.html>. Consider provisioning more memory on this node or reducing task parallelism by requesting more CPUs per task. To adjust the kill threshold, set the environment variable `RAY_memory_usage_threshold` when starting Ray. To disable worker killing, set the environment variable `RAY_memory_monitor_refresh_ms` to zero.Colin Ho
02/12/2025, 1:33 AMdf.explain(True) ) ?Neil Wadhvana
02/12/2025, 1:33 AMNeil Wadhvana
02/12/2025, 1:34 AMColin Ho
02/12/2025, 1:34 AMColin Ho
02/12/2025, 1:34 AMNeil Wadhvana
02/12/2025, 1:34 AMRay is the one that reschedules itIs it triggering the
retry or restart on the actor?Colin Ho
02/12/2025, 1:35 AMNeil Wadhvana
02/12/2025, 1:35 AMColin Ho
02/12/2025, 1:35 AMNeil Wadhvana
02/12/2025, 1:36 AMNeil Wadhvana
02/12/2025, 1:38 AMjay
02/12/2025, 1:38 AMNeil Wadhvana
02/12/2025, 1:38 AMjay
02/12/2025, 1:39 AMThe native runner does not partition at all.And yes the native runner could actually be better here because it does stuff in small morsel chunks. I see your point wrt needing to run this in a Ray cluster eventually though. Would love to take a look at the plan once you have it
Neil Wadhvana
02/12/2025, 1:40 AMNeil Wadhvana
02/12/2025, 1:42 AMNeil Wadhvana
02/12/2025, 1:42 AMjay
02/12/2025, 1:45 AMThe ec2 autoscaler already helped us deliver this data processing jobUsing Daft in native mode? Or what do you mean by this
Neil Wadhvana
02/12/2025, 1:47 AM