mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-05 21:38:12 +08:00
Change imports for cluster and tasks
This commit is contained in:
+1
-2
@@ -21,8 +21,7 @@ from django.utils.translation import ugettext_lazy as _
|
|||||||
from django import db
|
from django import db
|
||||||
|
|
||||||
# Local
|
# Local
|
||||||
import tasks
|
from django_q import tasks
|
||||||
|
|
||||||
from django_q.compat import range
|
from django_q.compat import range
|
||||||
from django_q.conf import Conf, logger, psutil, get_ppid, error_reporter, rollbar
|
from django_q.conf import Conf, logger, psutil, get_ppid, error_reporter, rollbar
|
||||||
from django_q.models import Task, Success, Schedule
|
from django_q.models import Task, Success, Schedule
|
||||||
|
|||||||
+35
-34
@@ -1,4 +1,6 @@
|
|||||||
"""Provides task functionality."""
|
"""Provides task functionality."""
|
||||||
|
# Standard
|
||||||
|
from time import sleep, time
|
||||||
from multiprocessing import Value
|
from multiprocessing import Value
|
||||||
|
|
||||||
# django
|
# django
|
||||||
@@ -6,9 +8,8 @@ from django.db import IntegrityError
|
|||||||
from django.utils import timezone
|
from django.utils import timezone
|
||||||
|
|
||||||
# local
|
# local
|
||||||
import time
|
|
||||||
from django_q.signing import SignedPackage
|
from django_q.signing import SignedPackage
|
||||||
import cluster
|
from django_q import cluster
|
||||||
from django_q.conf import Conf, logger
|
from django_q.conf import Conf, logger
|
||||||
from django_q.models import Schedule, Task
|
from django_q.models import Schedule, Task
|
||||||
from django_q.humanhash import uuid
|
from django_q.humanhash import uuid
|
||||||
@@ -112,14 +113,14 @@ def result(task_id, wait=0, cached=Conf.CACHED):
|
|||||||
"""
|
"""
|
||||||
if cached:
|
if cached:
|
||||||
return result_cached(task_id, wait)
|
return result_cached(task_id, wait)
|
||||||
start = time.time()
|
start = time()
|
||||||
while True:
|
while True:
|
||||||
r = Task.get_result(task_id)
|
r = Task.get_result(task_id)
|
||||||
if r:
|
if r:
|
||||||
return r
|
return r
|
||||||
if (time.time() - start) * 1000 >= wait >= 0:
|
if (time() - start) * 1000 >= wait >= 0:
|
||||||
break
|
break
|
||||||
time.sleep(0.01)
|
sleep(0.01)
|
||||||
|
|
||||||
|
|
||||||
def result_cached(task_id, wait=0, broker=None):
|
def result_cached(task_id, wait=0, broker=None):
|
||||||
@@ -128,14 +129,14 @@ def result_cached(task_id, wait=0, broker=None):
|
|||||||
"""
|
"""
|
||||||
if not broker:
|
if not broker:
|
||||||
broker = get_broker()
|
broker = get_broker()
|
||||||
start = time.time()
|
start = time()
|
||||||
while True:
|
while True:
|
||||||
r = broker.cache.get('{}:{}'.format(broker.list_key, task_id))
|
r = broker.cache.get('{}:{}'.format(broker.list_key, task_id))
|
||||||
if r:
|
if r:
|
||||||
return SignedPackage.loads(r)['result']
|
return SignedPackage.loads(r)['result']
|
||||||
if (time.time() - start) * 1000 >= wait >= 0:
|
if (time() - start) * 1000 >= wait >= 0:
|
||||||
break
|
break
|
||||||
time.sleep(0.01)
|
sleep(0.01)
|
||||||
|
|
||||||
|
|
||||||
def result_group(group_id, failures=False, wait=0, count=None, cached=Conf.CACHED):
|
def result_group(group_id, failures=False, wait=0, count=None, cached=Conf.CACHED):
|
||||||
@@ -150,19 +151,19 @@ def result_group(group_id, failures=False, wait=0, count=None, cached=Conf.CACHE
|
|||||||
"""
|
"""
|
||||||
if cached:
|
if cached:
|
||||||
return result_group_cached(group_id, failures, wait, count)
|
return result_group_cached(group_id, failures, wait, count)
|
||||||
start = time.time()
|
start = time()
|
||||||
if count:
|
if count:
|
||||||
while True:
|
while True:
|
||||||
if count_group(group_id) == count or wait and (time.time() - start) * 1000 >= wait >= 0:
|
if count_group(group_id) == count or wait and (time() - start) * 1000 >= wait >= 0:
|
||||||
break
|
break
|
||||||
time.sleep(0.01)
|
sleep(0.01)
|
||||||
while True:
|
while True:
|
||||||
r = Task.get_result_group(group_id, failures)
|
r = Task.get_result_group(group_id, failures)
|
||||||
if r:
|
if r:
|
||||||
return r
|
return r
|
||||||
if (time.time() - start) * 1000 >= wait >= 0:
|
if (time() - start) * 1000 >= wait >= 0:
|
||||||
break
|
break
|
||||||
time.sleep(0.01)
|
sleep(0.01)
|
||||||
|
|
||||||
|
|
||||||
def result_group_cached(group_id, failures=False, wait=0, count=None, broker=None):
|
def result_group_cached(group_id, failures=False, wait=0, count=None, broker=None):
|
||||||
@@ -171,12 +172,12 @@ def result_group_cached(group_id, failures=False, wait=0, count=None, broker=Non
|
|||||||
"""
|
"""
|
||||||
if not broker:
|
if not broker:
|
||||||
broker = get_broker()
|
broker = get_broker()
|
||||||
start = time.time()
|
start = time()
|
||||||
if count:
|
if count:
|
||||||
while True:
|
while True:
|
||||||
if count_group_cached(group_id) == count or wait and (time.time() - start) * 1000 >= wait > 0:
|
if count_group_cached(group_id) == count or wait and (time() - start) * 1000 >= wait > 0:
|
||||||
break
|
break
|
||||||
time.sleep(0.01)
|
sleep(0.01)
|
||||||
while True:
|
while True:
|
||||||
group_list = broker.cache.get('{}:{}:keys'.format(broker.list_key, group_id))
|
group_list = broker.cache.get('{}:{}:keys'.format(broker.list_key, group_id))
|
||||||
if group_list:
|
if group_list:
|
||||||
@@ -186,9 +187,9 @@ def result_group_cached(group_id, failures=False, wait=0, count=None, broker=Non
|
|||||||
if task['success'] or failures:
|
if task['success'] or failures:
|
||||||
result_list.append(task['result'])
|
result_list.append(task['result'])
|
||||||
return result_list
|
return result_list
|
||||||
if (time.time() - start) * 1000 >= wait >= 0:
|
if (time() - start) * 1000 >= wait >= 0:
|
||||||
break
|
break
|
||||||
time.sleep(0.01)
|
sleep(0.01)
|
||||||
|
|
||||||
|
|
||||||
def fetch(task_id, wait=0, cached=Conf.CACHED):
|
def fetch(task_id, wait=0, cached=Conf.CACHED):
|
||||||
@@ -205,14 +206,14 @@ def fetch(task_id, wait=0, cached=Conf.CACHED):
|
|||||||
"""
|
"""
|
||||||
if cached:
|
if cached:
|
||||||
return fetch_cached(task_id, wait)
|
return fetch_cached(task_id, wait)
|
||||||
start = time.time()
|
start = time()
|
||||||
while True:
|
while True:
|
||||||
t = Task.get_task(task_id)
|
t = Task.get_task(task_id)
|
||||||
if t:
|
if t:
|
||||||
return t
|
return t
|
||||||
if (time.time() - start) * 1000 >= wait >= 0:
|
if (time() - start) * 1000 >= wait >= 0:
|
||||||
break
|
break
|
||||||
time.sleep(0.01)
|
sleep(0.01)
|
||||||
|
|
||||||
|
|
||||||
def fetch_cached(task_id, wait=0, broker=None):
|
def fetch_cached(task_id, wait=0, broker=None):
|
||||||
@@ -221,7 +222,7 @@ def fetch_cached(task_id, wait=0, broker=None):
|
|||||||
"""
|
"""
|
||||||
if not broker:
|
if not broker:
|
||||||
broker = get_broker()
|
broker = get_broker()
|
||||||
start = time.time()
|
start = time()
|
||||||
while True:
|
while True:
|
||||||
r = broker.cache.get('{}:{}'.format(broker.list_key, task_id))
|
r = broker.cache.get('{}:{}'.format(broker.list_key, task_id))
|
||||||
if r:
|
if r:
|
||||||
@@ -237,9 +238,9 @@ def fetch_cached(task_id, wait=0, broker=None):
|
|||||||
result=task['result'],
|
result=task['result'],
|
||||||
success=task['success'])
|
success=task['success'])
|
||||||
return t
|
return t
|
||||||
if (time.time() - start) * 1000 >= wait >= 0:
|
if (time() - start) * 1000 >= wait >= 0:
|
||||||
break
|
break
|
||||||
time.sleep(0.01)
|
sleep(0.01)
|
||||||
|
|
||||||
|
|
||||||
def fetch_group(group_id, failures=True, wait=0, count=None, cached=Conf.CACHED):
|
def fetch_group(group_id, failures=True, wait=0, count=None, cached=Conf.CACHED):
|
||||||
@@ -253,19 +254,19 @@ def fetch_group(group_id, failures=True, wait=0, count=None, cached=Conf.CACHED)
|
|||||||
"""
|
"""
|
||||||
if cached:
|
if cached:
|
||||||
return fetch_group_cached(group_id, failures, wait, count)
|
return fetch_group_cached(group_id, failures, wait, count)
|
||||||
start = time.time()
|
start = time()
|
||||||
if count:
|
if count:
|
||||||
while True:
|
while True:
|
||||||
if count_group(group_id) == count or wait and (time.time() - start) * 1000 >= wait >= 0:
|
if count_group(group_id) == count or wait and (time() - start) * 1000 >= wait >= 0:
|
||||||
break
|
break
|
||||||
time.sleep(0.01)
|
sleep(0.01)
|
||||||
while True:
|
while True:
|
||||||
r = Task.get_task_group(group_id, failures)
|
r = Task.get_task_group(group_id, failures)
|
||||||
if r:
|
if r:
|
||||||
return r
|
return r
|
||||||
if (time.time() - start) * 1000 >= wait >= 0:
|
if (time() - start) * 1000 >= wait >= 0:
|
||||||
break
|
break
|
||||||
time.sleep(0.01)
|
sleep(0.01)
|
||||||
|
|
||||||
|
|
||||||
def fetch_group_cached(group_id, failures=True, wait=0, count=None, broker=None):
|
def fetch_group_cached(group_id, failures=True, wait=0, count=None, broker=None):
|
||||||
@@ -274,12 +275,12 @@ def fetch_group_cached(group_id, failures=True, wait=0, count=None, broker=None)
|
|||||||
"""
|
"""
|
||||||
if not broker:
|
if not broker:
|
||||||
broker = get_broker()
|
broker = get_broker()
|
||||||
start = time.time()
|
start = time()
|
||||||
if count:
|
if count:
|
||||||
while True:
|
while True:
|
||||||
if count_group_cached(group_id) == count or wait and (time.time() - start) * 1000 >= wait >= 0:
|
if count_group_cached(group_id) == count or wait and (time() - start) * 1000 >= wait >= 0:
|
||||||
break
|
break
|
||||||
time.sleep(0.01)
|
sleep(0.01)
|
||||||
while True:
|
while True:
|
||||||
group_list = broker.cache.get('{}:{}:keys'.format(broker.list_key, group_id))
|
group_list = broker.cache.get('{}:{}:keys'.format(broker.list_key, group_id))
|
||||||
if group_list:
|
if group_list:
|
||||||
@@ -300,9 +301,9 @@ def fetch_group_cached(group_id, failures=True, wait=0, count=None, broker=None)
|
|||||||
success=task['success'])
|
success=task['success'])
|
||||||
task_list.append(t)
|
task_list.append(t)
|
||||||
return task_list
|
return task_list
|
||||||
if (time.time() - start) * 1000 >= wait >= 0:
|
if (time() - start) * 1000 >= wait >= 0:
|
||||||
break
|
break
|
||||||
time.sleep(0.01)
|
sleep(0.01)
|
||||||
|
|
||||||
|
|
||||||
def count_group(group_id, failures=False, cached=Conf.CACHED):
|
def count_group(group_id, failures=False, cached=Conf.CACHED):
|
||||||
|
|||||||
Reference in New Issue
Block a user