mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-07 03:58:11 +08:00
Revert "Merge pull request #15 from django-q2/error-on-timeout"
This reverts commitd97dca5bcb, reversing changes made to01cb652071.
This commit is contained in:
@@ -458,29 +458,6 @@ def monitor(result_queue: Queue, broker: Broker = None):
|
||||
logger.info(_("%(name)s stopped monitoring results") % {"name": proc_name})
|
||||
|
||||
|
||||
def _check_task_timed_out(key, task: dict):
|
||||
result = None
|
||||
broker = get_broker()
|
||||
cache = broker.cache
|
||||
working_set = cache.get(key) or set()
|
||||
if task["id"] in working_set:
|
||||
# the previous worker has timedout and wasn't given chance to clear
|
||||
raise Exception(f"Task Timed-out: {task}.")
|
||||
else:
|
||||
working_set.add(task['id'])
|
||||
cache.set(key, working_set, timeout=Conf.RETRY * 3)
|
||||
return result
|
||||
|
||||
|
||||
def _clear_task_timeout_cache(key, task):
|
||||
broker = get_broker()
|
||||
cache = broker.cache
|
||||
working_set = cache.get(key) or set()
|
||||
if task["id"] in working_set:
|
||||
working_set.remove(task['id'])
|
||||
cache.set(key, working_set, timeout=Conf.RETRY * 3)
|
||||
|
||||
|
||||
def worker(
|
||||
task_queue: Queue, result_queue: Queue, timer: Value, timeout: int = Conf.TIMEOUT
|
||||
):
|
||||
@@ -502,8 +479,6 @@ def worker(
|
||||
task_count = 0
|
||||
if timeout is None:
|
||||
timeout = -1
|
||||
|
||||
working_tasks_key = "DJANGO-Q-WORKING-TASKS"
|
||||
# Start reading the task queue
|
||||
for task in iter(task_queue.get, "STOP"):
|
||||
result = None
|
||||
@@ -543,11 +518,7 @@ def worker(
|
||||
pre_execute.send(sender="django_q", func=f, task=task)
|
||||
# execute the payload
|
||||
timer.value = timer_value # Busy
|
||||
|
||||
|
||||
try:
|
||||
if Conf.FAIL_ON_TIMEOUT:
|
||||
_check_task_timed_out(working_tasks_key, task)
|
||||
if f is None:
|
||||
# raise a meaningfull error if task["func"] is not a valid function
|
||||
raise ValueError(f"Function {task['func']} is not defined")
|
||||
@@ -559,9 +530,6 @@ def worker(
|
||||
error_reporter.report()
|
||||
if task.get("sync", False):
|
||||
raise
|
||||
if Conf.FAIL_ON_TIMEOUT:
|
||||
_clear_task_timeout_cache(working_tasks_key, task)
|
||||
|
||||
with timer.get_lock():
|
||||
# Process result
|
||||
task["result"] = result[0]
|
||||
|
||||
@@ -133,9 +133,6 @@ class Conf:
|
||||
# Number of seconds to wait for a worker to finish.
|
||||
TIMEOUT = conf.get("timeout", None)
|
||||
|
||||
# Whether to fail the task when it times-out.
|
||||
FAIL_ON_TIMEOUT = conf.get("fail_on_timeout", False)
|
||||
|
||||
# Whether to acknowledge unsuccessful tasks.
|
||||
# This causes failed tasks to be considered delivered, thereby removing them from
|
||||
# the task queue. Defaults to False.
|
||||
|
||||
Reference in New Issue
Block a user