mirror of
https://github.com/django-q2/django-q2.git
synced 2026-10-08 03:38:11 +08:00
Adds max_rss memory limit for recycle
This commit is contained in:
+18
-16
@@ -1,6 +1,7 @@
|
|||||||
# Standard
|
# Standard
|
||||||
import ast
|
import ast
|
||||||
import importlib
|
import importlib
|
||||||
|
import resource
|
||||||
import signal
|
import signal
|
||||||
import socket
|
import socket
|
||||||
import traceback
|
import traceback
|
||||||
@@ -10,6 +11,7 @@ from time import sleep
|
|||||||
|
|
||||||
# External
|
# External
|
||||||
import arrow
|
import arrow
|
||||||
|
|
||||||
# Django
|
# Django
|
||||||
from django import db
|
from django import db
|
||||||
from django.conf import settings
|
from django.conf import settings
|
||||||
@@ -102,10 +104,10 @@ class Cluster:
|
|||||||
@property
|
@property
|
||||||
def is_stopping(self) -> bool:
|
def is_stopping(self) -> bool:
|
||||||
return (
|
return (
|
||||||
self.stop_event
|
self.stop_event
|
||||||
and self.start_event
|
and self.start_event
|
||||||
and self.start_event.is_set()
|
and self.start_event.is_set()
|
||||||
and self.stop_event.is_set()
|
and self.stop_event.is_set()
|
||||||
)
|
)
|
||||||
|
|
||||||
@property
|
@property
|
||||||
@@ -115,13 +117,13 @@ class Cluster:
|
|||||||
|
|
||||||
class Sentinel:
|
class Sentinel:
|
||||||
def __init__(
|
def __init__(
|
||||||
self,
|
self,
|
||||||
stop_event,
|
stop_event,
|
||||||
start_event,
|
start_event,
|
||||||
cluster_id,
|
cluster_id,
|
||||||
broker=None,
|
broker=None,
|
||||||
timeout=Conf.TIMEOUT,
|
timeout=Conf.TIMEOUT,
|
||||||
start=True,
|
start=True,
|
||||||
):
|
):
|
||||||
# Make sure we catch signals for the pool
|
# Make sure we catch signals for the pool
|
||||||
signal.signal(signal.SIGINT, signal.SIG_IGN)
|
signal.signal(signal.SIGINT, signal.SIG_IGN)
|
||||||
@@ -376,7 +378,7 @@ def monitor(result_queue: Queue, broker: Broker = None):
|
|||||||
|
|
||||||
|
|
||||||
def worker(
|
def worker(
|
||||||
task_queue: Queue, result_queue: Queue, timer: Value, timeout: int = Conf.TIMEOUT
|
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
|
||||||
@@ -433,7 +435,7 @@ def worker(
|
|||||||
result_queue.put(task)
|
result_queue.put(task)
|
||||||
timer.value = -1 # Idle
|
timer.value = -1 # Idle
|
||||||
# Recycle
|
# Recycle
|
||||||
if task_count == Conf.RECYCLE:
|
if task_count == Conf.RECYCLE or (Conf.MAX_RSS and resource.getrusage(resource.RUSAGE_SELF).ru_maxrss >= Conf.MAX_RSS):
|
||||||
timer.value = -2 # Recycled
|
timer.value = -2 # Recycled
|
||||||
break
|
break
|
||||||
logger.info(_(f"{name} stopped doing work"))
|
logger.info(_(f"{name} stopped doing work"))
|
||||||
@@ -551,9 +553,9 @@ def scheduler(broker: Broker = None):
|
|||||||
try:
|
try:
|
||||||
with db.transaction.atomic(using=Schedule.objects.db):
|
with db.transaction.atomic(using=Schedule.objects.db):
|
||||||
for s in (
|
for s in (
|
||||||
Schedule.objects.select_for_update()
|
Schedule.objects.select_for_update()
|
||||||
.exclude(repeats=0)
|
.exclude(repeats=0)
|
||||||
.filter(next_run__lt=timezone.now())
|
.filter(next_run__lt=timezone.now())
|
||||||
):
|
):
|
||||||
args = ()
|
args = ()
|
||||||
kwargs = {}
|
kwargs = {}
|
||||||
|
|||||||
@@ -104,6 +104,9 @@ class Conf:
|
|||||||
# Number of tasks each worker can handle before it gets recycled. Useful for releasing memory
|
# Number of tasks each worker can handle before it gets recycled. Useful for releasing memory
|
||||||
RECYCLE = conf.get("recycle", 500)
|
RECYCLE = conf.get("recycle", 500)
|
||||||
|
|
||||||
|
# The maximum resident set size in kilobytes before a worker will recycle. Useful for limiting memory usage.
|
||||||
|
MAX_RSS = conf.get("max_rss", None)
|
||||||
|
|
||||||
# Number of seconds to wait for a worker to finish.
|
# Number of seconds to wait for a worker to finish.
|
||||||
TIMEOUT = conf.get("timeout", None)
|
TIMEOUT = conf.get("timeout", None)
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user