mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-05 17:28:12 +08:00
+5
-1
@@ -34,6 +34,7 @@ from django_q.conf import Conf, logger, psutil, get_ppid, rollbar
|
|||||||
from django_q.models import Task, Success, Schedule
|
from django_q.models import Task, Success, Schedule
|
||||||
from django_q.status import Stat, Status
|
from django_q.status import Stat, Status
|
||||||
from django_q.brokers import get_broker
|
from django_q.brokers import get_broker
|
||||||
|
from django_q.signals import pre_execute
|
||||||
|
|
||||||
|
|
||||||
class Cluster(object):
|
class Cluster(object):
|
||||||
@@ -373,8 +374,11 @@ def worker(task_queue, result_queue, timer, timeout=Conf.TIMEOUT):
|
|||||||
# We're still going
|
# We're still going
|
||||||
if not result:
|
if not result:
|
||||||
db.close_old_connections()
|
db.close_old_connections()
|
||||||
|
timer_value = task['kwargs'].pop('timeout', timeout or 0)
|
||||||
|
# signal execution
|
||||||
|
pre_execute.send(sender="django_q", func=f, task=task)
|
||||||
# execute the payload
|
# execute the payload
|
||||||
timer.value = task['kwargs'].pop('timeout', timeout or 0) # Busy
|
timer.value = timer_value # Busy
|
||||||
try:
|
try:
|
||||||
res = f(*task['args'], **task['kwargs'])
|
res = f(*task['args'], **task['kwargs'])
|
||||||
result = (res, True)
|
result = (res, True)
|
||||||
|
|||||||
+5
-1
@@ -1,7 +1,7 @@
|
|||||||
import importlib
|
import importlib
|
||||||
|
|
||||||
from django.db.models.signals import post_save
|
from django.db.models.signals import post_save
|
||||||
from django.dispatch import receiver
|
from django.dispatch import receiver, Signal
|
||||||
from django.utils.translation import ugettext_lazy as _
|
from django.utils.translation import ugettext_lazy as _
|
||||||
|
|
||||||
from django_q.conf import logger
|
from django_q.conf import logger
|
||||||
@@ -24,3 +24,7 @@ def call_hook(sender, instance, **kwargs):
|
|||||||
f(instance)
|
f(instance)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(_('return hook {} failed on [{}] because {}').format(instance.hook, instance.name, e))
|
logger.error(_('return hook {} failed on [{}] because {}').format(instance.hook, instance.name, e))
|
||||||
|
|
||||||
|
|
||||||
|
pre_enqueue = Signal(providing_args=["task"])
|
||||||
|
pre_execute = Signal(providing_args=["func", "task"])
|
||||||
|
|||||||
@@ -13,6 +13,7 @@ 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
|
||||||
from django_q.brokers import get_broker
|
from django_q.brokers import get_broker
|
||||||
|
from django_q.signals import pre_enqueue
|
||||||
|
|
||||||
|
|
||||||
def async(func, *args, **kwargs):
|
def async(func, *args, **kwargs):
|
||||||
@@ -43,6 +44,8 @@ def async(func, *args, **kwargs):
|
|||||||
# finalize
|
# finalize
|
||||||
task['kwargs'] = keywords
|
task['kwargs'] = keywords
|
||||||
task['started'] = timezone.now()
|
task['started'] = timezone.now()
|
||||||
|
# signal it
|
||||||
|
pre_enqueue.send(sender="django_q", task=task)
|
||||||
# sign it
|
# sign it
|
||||||
pack = signing.SignedPackage.dumps(task)
|
pack = signing.SignedPackage.dumps(task)
|
||||||
if task.get('sync', False):
|
if task.get('sync', False):
|
||||||
|
|||||||
@@ -42,6 +42,7 @@ Contents:
|
|||||||
Cluster <cluster>
|
Cluster <cluster>
|
||||||
Monitor <monitor>
|
Monitor <monitor>
|
||||||
Admin <admin>
|
Admin <admin>
|
||||||
|
Signals <signals>
|
||||||
Architecture <architecture>
|
Architecture <architecture>
|
||||||
Examples <examples>
|
Examples <examples>
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,43 @@
|
|||||||
|
Signals
|
||||||
|
=======
|
||||||
|
.. py:currentmodule:: django_q
|
||||||
|
|
||||||
|
Available signals
|
||||||
|
-----------------
|
||||||
|
|
||||||
|
Django Q emits the following signals during its lifecycle.
|
||||||
|
|
||||||
|
Before enqueuing a task
|
||||||
|
"""""""""""""""""""""""
|
||||||
|
|
||||||
|
The ``django_q.signals.pre_enqueue`` signal is emitted before a task is
|
||||||
|
enqueued. The task dictionary is given as the ``task`` argument.
|
||||||
|
|
||||||
|
Before executing a task
|
||||||
|
"""""""""""""""""""""""
|
||||||
|
|
||||||
|
The ``django_q.signals.pre_execute`` signal is emitted before a task is
|
||||||
|
executed by a worker. This signal provides two arguments:
|
||||||
|
|
||||||
|
- ``task``: the task dictionary.
|
||||||
|
- ``func``: the actual function that will be executed. If the task was created
|
||||||
|
with a function path, this argument will be the callable function
|
||||||
|
nonetheless.
|
||||||
|
|
||||||
|
Subscribing to a signal
|
||||||
|
-----------------------
|
||||||
|
|
||||||
|
Connecting to a Django Q signal is done in the same manner as any other Django
|
||||||
|
signal::
|
||||||
|
|
||||||
|
from django.dispatch import receiver
|
||||||
|
from django_q.signals import pre_enqueue, pre_execute
|
||||||
|
|
||||||
|
@receiver(pre_enqueue)
|
||||||
|
def my_pre_enqueue_callback(sender, task, **kwargs):
|
||||||
|
print("Task {} will be enqueued".format(task["name"]))
|
||||||
|
|
||||||
|
@receiver(pre_execute)
|
||||||
|
def my_pre_execute_callback(sender, func, task, **kwargs):
|
||||||
|
print("Task {} will be executed by calling {}".format(
|
||||||
|
task["name"], func))
|
||||||
Reference in New Issue
Block a user