Re: Ray OOM retries Ray seems to be quite aggress...
# daft-dev
j
Re: Ray OOM retries Ray seems to be quite aggressive about OOM retries and seems to usually retry enough till it completes runs — i.e. all the Partitions get correctly created. However it seems that we still encounter errors further downstream after the run completes… Interestingly this is when we attempt to retrieve partition metadatas for displaying the length of the df: 1. Because Partition metadatas are small-enough, Ray keeps them in worker memory 2. When a worker subsequently OOMs (while working on some other task), this partition metadata is lost to the wind… 3. Later on when we attempt to retrieve it (for visualization of total length), we get a WorkerCrashedError:
Copy code
Traceback (most recent call last):
  File "/tmp/ray/session_2024-12-11_06-29-26_666316_2025/runtime_resources/working_dir_files/_ray_pkg_3abbaf3dbef03d3e/big_task_heap_usage.py", line 70, in <module>
    df.collect()
...
  File "/home/ubuntu/.venv/lib/python3.9/site-packages/daft/dataframe/dataframe.py", line 2531, in collect
    dataframe_len = len(self._result)
  File "/home/ubuntu/.venv/lib/python3.9/site-packages/daft/runners/ray_runner.py", line 329, in __len__
    return sum(result.metadata().num_rows for result in self._results.values())
  File "/home/ubuntu/.venv/lib/python3.9/site-packages/daft/runners/ray_runner.py", line 329, in <genexpr>
    return sum(result.metadata().num_rows for result in self._results.values())
...
ray.exceptions.WorkerCrashedError: The worker died unexpectedly while executing this task. Check python-core-worker-*.log files for more information.
Example run: https://github.com/Eventual-Inc/Daft/actions/runs/12270896011/job/34236863613
yubikey intensifies 1
👀 1
d
Does it make a difference if there was 1 retry vs n retries?
j
Corrected my premature message 😛
😅 1
d
sounds like putting the partition metadata in the ray object store instead of worker memory could help with this?
j
Sometimes it still works though… I’m guessing Ray has its own redundancy (stores multiple copies in different workers), but I’m guessing this is happening if the OOMs kill enough of those workers
Perhaps the best solution here is our own metadata/metrics accounting, outside of Ray
nods 1
Ray Dataset for example does run a StatsActor, which could probably help here too since the Actor runs as a separate process.