-
Notifications
You must be signed in to change notification settings - Fork 2
/
daskmanager.py
49 lines (37 loc) · 1.5 KB
/
daskmanager.py
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
import logging
import traceback
from dask.distributed import Client, Future
from django.conf import settings
from daskmanager.models import DaskTask
logger = logging.getLogger(__name__)
class Singleton(type):
_instances = {}
def __call__(cls, *args, **kwargs):
if cls not in cls._instances:
cls._instances[cls] = super(Singleton, cls).__call__(*args, **kwargs)
return cls._instances[cls]
class DaskManager(metaclass=Singleton):
def __init__(self):
self.client = Client(f'{settings.DASK_SCHEDULER_HOST}:{settings.DASK_SCHEDULER_PORT}')
def compute(self, graph):
future = self.client.compute(graph)
future.add_done_callback(self.task_complete)
dask_task = DaskTask.objects.create(task_key=future.key)
return dask_task
def get_future_status(self, task_key):
return Future(key=task_key, client=self.client).status
@staticmethod
def task_complete(future):
task = DaskTask.objects.get(pk=future.key)
if future.status == 'finished':
task.status = future.status
task.result = future.result()
task.save()
elif future.status == 'error':
task.status = future.status
task.result = traceback.extract_tb(future.traceback()) + [future.exception()]
task.save()
# will cause exception to be thrown here
future.result()
else:
logger.error('Task completed with unhandled status: ' + future.status)