Hey guys! I'm tuning a multi-TB Multi-modal data t...
# daft-dev
n
Hey guys! I'm tuning a multi-TB Multi-modal data transformation pipeline which reads from unstructured source data and uses Daft to structure it. One of the requirements is to unpack data from sample-level binary files (ie. a npz file) and extract its content to separate, smaller, specifically- conditioned locations. This job OOMs on one of the largest AWS CPU instances using the Ray backend. When tuning this job, I tried three extremes (even knowing the limitations): 1. No partitioning - let daft do things a. Outcome: early-job OOM. Fails hard, fails fast. 2. Partitioning according to the docs with 2*96 (96 CPU cores) a. Outcome: early-job OOM since a single partition is too big to fit in memory 3. Partitioning to the number of samples (ie.
df.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?
👋 1
Also, I may just be overthinking this - when the ray kills raylets scheduled by the daft runner due to OOMs, does the daft runner simply re-schedule them to run again? This is from the ray docs, but I can't tell which of these cases the daft runner is triggering / running into, and what retry policy is used: https://docs.ray.io/en/latest/ray-core/scheduling/ray-oom-prevention.html#retry-policy
Example log output:
Copy code
(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.
c
Are you running on a single node? Also if possible, could you share the plan (
df.explain(True)
) ?
n
yep - will dump as soon as the execution stops (it's still running)
It's also in a notebook ... would this make any runtime difference ...? I don't think so, right?
c
Also, Daft schedules tasks to Ray. If the task gets killed (like in the log), Ray is the one that reschedules it
👍 1
Notebook should make no difference
n
Ray is the one that reschedules it
Is it triggering the
retry
or
restart
on the actor?
c
Also, is the default native runner not working for you?
n
(ie. can we safely assume ray re-scheduled the actor and will do so infinitely, as documented)
n
> Also, is the default native runner not working for you? The native runner does not partition at all. Further, we will be scaling this job across a ray cluster, not just a single node in the future. Future = 2 ish weeks from now once EKS is up and running. We've already tried this using the ray autoscaler and it works okay. Just trying to figure out how to debug such issues and interpret error messages when working on a single node.
Does daft retry indefinitely? Or is that a parameter we can tune?
j
Retries are handled by Ray atm which I believe retries indefinitely
❤️ 1
n
Amazing - we will simply let it run until something dies in this case. Not the best solution, but it works
j
The 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
👍 1
n
Np - will share once the notebook completes
🙌 1
The ec2 autoscaler already helped us deliver this data processing job using c7gn (arm) instances with 10GB/s saturated directly from S3 on each instance. All with a simple cluster config change.
The difference here is that we're limited to local memory I believe
j
The ec2 autoscaler already helped us deliver this data processing job
Using Daft in native mode? Or what do you mean by this
n
daft + ray runner + ray autoscaling
🙌 1