Using Ray in Python, I am doing a numerical simulation on 1000 elements running in parallel (on a PC with 8 cores). Sometimes the simulation gets stuck, and I would like to cancel the tasks that run for more than 90 seconds with ray.cancel(), so the rest of the elements in the queue can be processed. I can use ray.wait() to check if any simulations are finished and then process them, and I understand that best practice is to use a loop with ray.wait() instead of ray.get(). But among the unfinished tasks, I don't know how to distinguish between long-running tasks and pending tasks. As a workaround, I am thinking of using ray.get(elem, timeout = 90) in a for loop, and then cancelling the tasks that don't finish within the timeout, like this:
import ray
from ray.exceptions import GetTimeoutError
ray.init()
# Create jobs in element_list
# ...
# Collect finished jobs
finished = []
cancelled = []
for elem in element_list:
try:
finished.append(ray.get(elem,timeout=90))
except GetTimeoutError:
ray.cancel(elem,force=True)
cancelled.append(elem)
Is this the best way to do a thing like this? Or is it possible to get a list of the executing tasks and their running times? Or would it be possible to have a timer internally in the task that would return after 90 seconds without a solution? I am using a differential equation solver (scipy.integrate.solve_ivp), so it's not so easy to do a timeout check internally in the task.
Edit: I think this thread is discussing the same problem. And here is a feature request for the same thing. I will look into them.