mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-07 13:08:11 +08:00
Chore: Flake8, isort and Black (#40)
This commit is contained in:
+135
-44
@@ -74,7 +74,7 @@ class Cluster:
|
||||
),
|
||||
)
|
||||
self.sentinel.start()
|
||||
logger.info(_("Q Cluster %(name)s starting.") % {'name': self.name})
|
||||
logger.info(_("Q Cluster %(name)s starting.") % {"name": self.name})
|
||||
while not self.start_event.is_set():
|
||||
sleep(0.1)
|
||||
return self.pid
|
||||
@@ -82,19 +82,21 @@ class Cluster:
|
||||
def stop(self) -> bool:
|
||||
if not self.sentinel.is_alive():
|
||||
return False
|
||||
logger.info(_("Q Cluster %(name)s stopping.") % {'name': self.name})
|
||||
logger.info(_("Q Cluster %(name)s stopping.") % {"name": self.name})
|
||||
self.stop_event.set()
|
||||
self.sentinel.join()
|
||||
logger.info(_("Q Cluster %(name)s has stopped.") % {'name': self.name})
|
||||
logger.info(_("Q Cluster %(name)s has stopped.") % {"name": self.name})
|
||||
self.start_event = None
|
||||
self.stop_event = None
|
||||
return True
|
||||
|
||||
def sig_handler(self, signum, frame):
|
||||
logger.debug(
|
||||
_(
|
||||
'%(name)s got signal %(signal)s'
|
||||
) % {'name': current_process().name, 'signal': Conf.SIGNAL_NAMES.get(signum, "UNKNOWN")}
|
||||
_("%(name)s got signal %(signal)s")
|
||||
% {
|
||||
"name": current_process().name,
|
||||
"signal": Conf.SIGNAL_NAMES.get(signum, "UNKNOWN"),
|
||||
}
|
||||
)
|
||||
self.stop()
|
||||
|
||||
@@ -216,21 +218,34 @@ class Sentinel:
|
||||
db.connections.close_all()
|
||||
if process == self.monitor:
|
||||
self.monitor = self.spawn_monitor()
|
||||
logger.error(_("reincarnated monitor %(name)s after sudden death") % {'name': process.name})
|
||||
logger.error(
|
||||
_("reincarnated monitor %(name)s after sudden death")
|
||||
% {"name": process.name}
|
||||
)
|
||||
elif process == self.pusher:
|
||||
self.pusher = self.spawn_pusher()
|
||||
logger.error(_("reincarnated pusher %(name)s after sudden death") % {'name': process.name})
|
||||
logger.error(
|
||||
_("reincarnated pusher %(name)s after sudden death")
|
||||
% {"name": process.name}
|
||||
)
|
||||
else:
|
||||
self.pool.remove(process)
|
||||
self.spawn_worker()
|
||||
if process.timer.value == 0:
|
||||
# only need to terminate on timeout, otherwise we risk destabilizing the queues
|
||||
# only need to terminate on timeout, otherwise we risk destabilizing
|
||||
# the queues
|
||||
process.terminate()
|
||||
logger.warning(_("reincarnated worker %(name)s after timeout") % {'name': process.name})
|
||||
logger.warning(
|
||||
_("reincarnated worker %(name)s after timeout")
|
||||
% {"name": process.name}
|
||||
)
|
||||
elif int(process.timer.value) == -2:
|
||||
logger.info(_("recycled worker %(name)s") % {'name': process.name})
|
||||
logger.info(_("recycled worker %(name)s") % {"name": process.name})
|
||||
else:
|
||||
logger.error(_("reincarnated worker %(name)s after death") % {'name': process.name})
|
||||
logger.error(
|
||||
_("reincarnated worker %(name)s after death")
|
||||
% {"name": process.name}
|
||||
)
|
||||
|
||||
self.reincarnations += 1
|
||||
|
||||
@@ -252,13 +267,18 @@ class Sentinel:
|
||||
|
||||
def guard(self):
|
||||
logger.info(
|
||||
_(
|
||||
"%(name)s guarding cluster %(cluster_name)s"
|
||||
) % {'name': current_process().name, 'cluster_name': humanize(self.cluster_id.hex)}
|
||||
_("%(name)s guarding cluster %(cluster_name)s")
|
||||
% {
|
||||
"name": current_process().name,
|
||||
"cluster_name": humanize(self.cluster_id.hex),
|
||||
}
|
||||
)
|
||||
self.start_event.set()
|
||||
Stat(self).save()
|
||||
logger.info(_("Q Cluster %(cluster_name)s running.") % {'cluster_name': humanize(self.cluster_id.hex)})
|
||||
logger.info(
|
||||
_("Q Cluster %(cluster_name)s running.")
|
||||
% {"cluster_name": humanize(self.cluster_id.hex)}
|
||||
)
|
||||
counter = 0
|
||||
cycle = Conf.GUARD_CYCLE # guard loop sleep in seconds
|
||||
# Guard loop. Runs at least once
|
||||
@@ -292,7 +312,7 @@ class Sentinel:
|
||||
def stop(self):
|
||||
Stat(self).save()
|
||||
name = current_process().name
|
||||
logger.info(_("%(name)s stopping cluster processes") % {'name': name})
|
||||
logger.info(_("%(name)s stopping cluster processes") % {"name": name})
|
||||
# Stopping pusher
|
||||
self.event_out.set()
|
||||
# Wait for it to stop
|
||||
@@ -317,7 +337,7 @@ class Sentinel:
|
||||
self.result_queue.close()
|
||||
# Wait for the result queue to empty
|
||||
self.result_queue.join_thread()
|
||||
logger.info(_("%(name)s waiting for the monitor.") % {'name': name})
|
||||
logger.info(_("%(name)s waiting for the monitor.") % {"name": name})
|
||||
# Wait for everything to close or time out
|
||||
count = 0
|
||||
if not self.timeout:
|
||||
@@ -339,7 +359,10 @@ def pusher(task_queue: Queue, event: Event, broker: Broker = None):
|
||||
"""
|
||||
if not broker:
|
||||
broker = get_broker()
|
||||
logger.info(_("%(process_name)s pushing tasks at %(id)s") % {'process_name': current_process().name, 'id': current_process().pid})
|
||||
logger.info(
|
||||
_("%(process_name)s pushing tasks at %(id)s")
|
||||
% {"process_name": current_process().name, "id": current_process().pid}
|
||||
)
|
||||
while True:
|
||||
try:
|
||||
task_set = broker.dequeue()
|
||||
@@ -360,10 +383,12 @@ def pusher(task_queue: Queue, event: Event, broker: Broker = None):
|
||||
continue
|
||||
task["ack_id"] = ack_id
|
||||
task_queue.put(task)
|
||||
logger.debug(_("queueing from %(list_key)s") % {'list_key': broker.list_key})
|
||||
logger.debug(
|
||||
_("queueing from %(list_key)s") % {"list_key": broker.list_key}
|
||||
)
|
||||
if event.is_set():
|
||||
break
|
||||
logger.info(_("%(name)s stopped pushing tasks") % {'name': current_process().name})
|
||||
logger.info(_("%(name)s stopped pushing tasks") % {"name": current_process().name})
|
||||
|
||||
|
||||
def monitor(result_queue: Queue, broker: Broker = None):
|
||||
@@ -375,7 +400,9 @@ def monitor(result_queue: Queue, broker: Broker = None):
|
||||
if not broker:
|
||||
broker = get_broker()
|
||||
name = current_process().name
|
||||
logger.info(_("%(name)s monitoring at %(id)s") % {'name': name, 'id': current_process().pid})
|
||||
logger.info(
|
||||
_("%(name)s monitoring at %(id)s") % {"name": name, "id": current_process().pid}
|
||||
)
|
||||
for task in iter(result_queue.get, "STOP"):
|
||||
# save the result
|
||||
if task.get("cached", False):
|
||||
@@ -389,28 +416,42 @@ def monitor(result_queue: Queue, broker: Broker = None):
|
||||
# signal execution done
|
||||
post_execute.send(sender="django_q", task=task)
|
||||
# log the result
|
||||
info_name = get_func_repr(task['func'])
|
||||
info_name = get_func_repr(task["func"])
|
||||
if task["success"]:
|
||||
# log success
|
||||
logger.info(_("Processed '%(info_name)s' (%(task_name)s)") % {'info_name': info_name, 'task_name': task['name']})
|
||||
logger.info(
|
||||
_("Processed '%(info_name)s' (%(task_name)s)")
|
||||
% {"info_name": info_name, "task_name": task["name"]}
|
||||
)
|
||||
else:
|
||||
# log failure
|
||||
logger.error(_("Failed '%(info_name)s' (%(task_name)s) - %(task_result)s") % {'info_name': info_name, 'task_name': task['name'], 'task_result': task['result']})
|
||||
logger.info(_("%(name)s stopped monitoring results") % {'name': name})
|
||||
logger.error(
|
||||
_("Failed '%(info_name)s' (%(task_name)s) - %(task_result)s")
|
||||
% {
|
||||
"info_name": info_name,
|
||||
"task_name": task["name"],
|
||||
"task_result": task["result"],
|
||||
}
|
||||
)
|
||||
logger.info(_("%(name)s stopped monitoring results") % {"name": name})
|
||||
|
||||
|
||||
def worker(
|
||||
task_queue: Queue, result_queue: Queue, timer: Value, timeout: int = Conf.TIMEOUT
|
||||
):
|
||||
"""
|
||||
Takes a task from the task queue, tries to execute it and puts the result back in the result queue
|
||||
Takes a task from the task queue, tries to execute it and puts the result back in
|
||||
the result queue
|
||||
:param timeout: number of seconds wait for a worker to finish.
|
||||
:type task_queue: multiprocessing.Queue
|
||||
:type result_queue: multiprocessing.Queue
|
||||
:type timer: multiprocessing.Value
|
||||
"""
|
||||
proc_name = current_process().name
|
||||
logger.info(_("%(proc_name)s ready for work at %(id)s") % {'proc_name': proc_name, 'id': current_process().pid})
|
||||
logger.info(
|
||||
_("%(proc_name)s ready for work at %(id)s")
|
||||
% {"proc_name": proc_name, "id": current_process().pid}
|
||||
)
|
||||
task_count = 0
|
||||
if timeout is None:
|
||||
timeout = -1
|
||||
@@ -422,7 +463,14 @@ def worker(
|
||||
# Get the function from the task
|
||||
func = task["func"]
|
||||
func_name = get_func_repr(func)
|
||||
logger.info(_("%(proc_name)s processing '%(func_name)s' (%(task_name)s)") % {'proc_name': proc_name, 'func_name': func_name, 'task_name': task['name']})
|
||||
logger.info(
|
||||
_("%(proc_name)s processing '%(func_name)s' (%(task_name)s)")
|
||||
% {
|
||||
"proc_name": proc_name,
|
||||
"func_name": func_name,
|
||||
"task_name": task["name"],
|
||||
}
|
||||
)
|
||||
f = task["func"]
|
||||
# if it's not an instance try to get it from the string
|
||||
if not callable(task["func"]):
|
||||
@@ -437,7 +485,14 @@ def worker(
|
||||
res = f(*task["args"], **task["kwargs"])
|
||||
result = (res, True)
|
||||
except Exception:
|
||||
result = (_("Could not process '%(func_name)s'. Check the location of the function and the args/kwargs.") % {'func_name': func_name}, False)
|
||||
result = (
|
||||
_(
|
||||
"Could not process '%(func_name)s'. Check the location of the "
|
||||
"function and the args/kwargs."
|
||||
)
|
||||
% {"func_name": func_name},
|
||||
False,
|
||||
)
|
||||
if error_reporter:
|
||||
error_reporter.report()
|
||||
if task.get("sync", False):
|
||||
@@ -453,7 +508,8 @@ def worker(
|
||||
if task_count == Conf.RECYCLE or rss_check():
|
||||
timer.value = -2 # Recycled
|
||||
break
|
||||
logger.info(_("%(proc_name)s stopped doing work") % {'proc_name': proc_name})
|
||||
logger.info(_("%(proc_name)s stopped doing work") % {"proc_name": proc_name})
|
||||
|
||||
|
||||
def save_task(task, broker: Broker):
|
||||
"""
|
||||
@@ -478,16 +534,27 @@ def save_task(task, broker: Broker):
|
||||
|
||||
try:
|
||||
filters = {}
|
||||
if Conf.SAVE_LIMIT_PER and Conf.SAVE_LIMIT_PER in {"group", "name", "func"} and Conf.SAVE_LIMIT_PER in task:
|
||||
if (
|
||||
Conf.SAVE_LIMIT_PER
|
||||
and Conf.SAVE_LIMIT_PER in {"group", "name", "func"}
|
||||
and Conf.SAVE_LIMIT_PER in task
|
||||
):
|
||||
value = task[Conf.SAVE_LIMIT_PER]
|
||||
if Conf.SAVE_LIMIT_PER == "func":
|
||||
value = get_func_repr(value)
|
||||
filters[Conf.SAVE_LIMIT_PER] = value
|
||||
|
||||
database_to_use = {"using": Conf.ORM if Conf.ORM else Schedule.objects.db} if not Conf.HAS_REPLICA else {}
|
||||
database_to_use = (
|
||||
{"using": Conf.ORM if Conf.ORM else Schedule.objects.db}
|
||||
if not Conf.HAS_REPLICA
|
||||
else {}
|
||||
)
|
||||
with db.transaction.atomic(**database_to_use):
|
||||
last = Success.objects.filter(**filters).select_for_update().last()
|
||||
if task["success"] and 0 < Conf.SAVE_LIMIT <= Success.objects.filter(**filters).count():
|
||||
if (
|
||||
task["success"]
|
||||
and 0 < Conf.SAVE_LIMIT <= Success.objects.filter(**filters).count()
|
||||
):
|
||||
last.delete()
|
||||
|
||||
# check if this task has previous results
|
||||
@@ -588,7 +655,11 @@ def scheduler(broker: Broker = None):
|
||||
broker = get_broker()
|
||||
close_old_django_connections()
|
||||
try:
|
||||
database_to_use = {"using": Conf.ORM if Conf.ORM else Schedule.objects.db} if not Conf.HAS_REPLICA else {}
|
||||
database_to_use = (
|
||||
{"using": Conf.ORM if Conf.ORM else Schedule.objects.db}
|
||||
if not Conf.HAS_REPLICA
|
||||
else {}
|
||||
)
|
||||
with db.transaction.atomic(**database_to_use):
|
||||
for s in (
|
||||
Schedule.objects.select_for_update()
|
||||
@@ -608,8 +679,13 @@ def scheduler(broker: Broker = None):
|
||||
except (SyntaxError, ValueError):
|
||||
# else use the kwargs syntax
|
||||
try:
|
||||
parsed_kwargs = ast.parse(f"f({s.kwargs})").body[0].value.keywords
|
||||
kwargs = {kwarg.arg: ast.literal_eval(kwarg.value) for kwarg in parsed_kwargs}
|
||||
parsed_kwargs = (
|
||||
ast.parse(f"f({s.kwargs})").body[0].value.keywords
|
||||
)
|
||||
kwargs = {
|
||||
kwarg.arg: ast.literal_eval(kwarg.value)
|
||||
for kwarg in parsed_kwargs
|
||||
}
|
||||
except (SyntaxError, ValueError):
|
||||
kwargs = {}
|
||||
if s.args:
|
||||
@@ -646,7 +722,8 @@ def scheduler(broker: Broker = None):
|
||||
if not croniter:
|
||||
raise ImportError(
|
||||
_(
|
||||
"Please install croniter to enable cron expressions"
|
||||
"Please install croniter to enable cron "
|
||||
"expressions"
|
||||
)
|
||||
)
|
||||
next_run = croniter(s.cron, localtime()).get_next(datetime)
|
||||
@@ -659,7 +736,8 @@ def scheduler(broker: Broker = None):
|
||||
scheduled_broker = broker
|
||||
try:
|
||||
scheduled_broker = get_broker(q_options["broker_name"])
|
||||
except: # invalid broker_name or non existing broker with broker_name
|
||||
except: # noqa: E722
|
||||
# invalid broker_name or non existing broker with broker_name
|
||||
pass
|
||||
q_options["broker"] = scheduled_broker
|
||||
q_options["group"] = q_options.get("group", s.name or s.id)
|
||||
@@ -669,14 +747,24 @@ def scheduler(broker: Broker = None):
|
||||
if not s.task:
|
||||
logger.error(
|
||||
_(
|
||||
"%(process_name)s failed to create a task from schedule [%(schedule)s]"
|
||||
) % {'process_name': current_process().name, 'schedule': s.name or s.id}
|
||||
"%(process_name)s failed to create a task from schedule "
|
||||
"[%(schedule)s]"
|
||||
)
|
||||
% {
|
||||
"process_name": current_process().name,
|
||||
"schedule": s.name or s.id,
|
||||
}
|
||||
)
|
||||
else:
|
||||
logger.info(
|
||||
_(
|
||||
"%(process_name)s created a task from schedule [%(schedule)s]"
|
||||
) % {'process_name': current_process().name, 'schedule': s.name or s.id}
|
||||
"%(process_name)s created a task from schedule "
|
||||
"[%(schedule)s]"
|
||||
)
|
||||
% {
|
||||
"process_name": current_process().name,
|
||||
"schedule": s.name or s.id,
|
||||
}
|
||||
)
|
||||
# default behavior is to delete a ONCE schedule
|
||||
if s.schedule_type == s.ONCE:
|
||||
@@ -741,7 +829,10 @@ def set_cpu_affinity(n: int, process_ids: list, actual: bool = not Conf.TESTING)
|
||||
p = psutil.Process(pid)
|
||||
if actual:
|
||||
p.cpu_affinity(affinity)
|
||||
logger.info(_("%(pid)s will use cpu %(affinity)s") % {'pid': pid, 'affinity': affinity})
|
||||
logger.info(
|
||||
_("%(pid)s will use cpu %(affinity)s")
|
||||
% {"pid": pid, "affinity": affinity}
|
||||
)
|
||||
|
||||
|
||||
def rss_check():
|
||||
|
||||
Reference in New Issue
Block a user