Skip to content

Commit 60a19f1

Browse files
committed
Better management of known queues
1 parent 5c2c663 commit 60a19f1

5 files changed

Lines changed: 58 additions & 14 deletions

File tree

mrq/basetasks/cleaning.py

Lines changed: 22 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -205,13 +205,33 @@ class CleanKnownQueues(Task):
205205
def run(self, params):
206206

207207
max_age = int(params.get("max_age") or (7 * 86400))
208+
pretend = bool(params.get("pretend"))
209+
check_mongo = bool(params.get("check_mongo"))
208210

209211
known_queues = Queue.redis_known_queues()
210212

213+
removed_queues = []
214+
215+
queues_from_config = Queue.all_known_from_config()
216+
217+
print "Found %s known queues & %s from config" % (len(known_queues), len(queues_from_config))
218+
211219
# Only clean queues older than N days
212220
time_threshold = time.time() - max_age
213221
for queue, time_last_used in known_queues.iteritems():
222+
if queue in queues_from_config:
223+
continue
214224
if time_last_used < time_threshold:
215225
q = Queue(queue, add_to_known_queues=False)
216-
if q.size() == 0:
217-
q.remove_from_known_queues()
226+
size = q.size()
227+
if check_mongo:
228+
size += connections.mongodb_jobs.mrq_jobs.count({"queue": queue})
229+
if size == 0:
230+
removed_queues.append(queue)
231+
print "Removing empty queue '%s' from known queues ..." % queue
232+
if not pretend:
233+
q.remove_from_known_queues()
234+
235+
print "Cleaned %s queues" % len(removed_queues)
236+
237+
return removed_queues

mrq/basetasks/utils.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -132,7 +132,7 @@ def perform_action(self, action, query, destination_queue):
132132

133133
# Between these two lines, jobs can become "lost" too.
134134

135-
Queue(destination_queue or queue).enqueue_job_ids(
135+
Queue(destination_queue or queue, add_to_known_queues=True).enqueue_job_ids(
136136
[str(x) for x in jobs_by_queue[queue]])
137137

138138
print stats

mrq/job.py

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -250,7 +250,7 @@ def requeue(self, queue=None, retry_count=0):
250250
queue = self.data["queue"]
251251

252252
from .queue import Queue
253-
queue_obj = Queue(queue)
253+
queue_obj = Queue(queue, add_to_known_queues=True)
254254

255255
self._save_status("queued", updates={
256256
"queue": queue,
@@ -567,7 +567,7 @@ def queue_raw_jobs(queue, params_list, **kwargs):
567567
""" Queue some jobs on a raw queue """
568568

569569
from .queue import Queue
570-
queue_obj = Queue(queue)
570+
queue_obj = Queue(queue, add_to_known_queues=True)
571571
queue_obj.enqueue_raw_jobs(params_list, **kwargs)
572572

573573

@@ -588,7 +588,7 @@ def queue_jobs(main_task_path, params_list, queue=None, batch_size=1000):
588588
queue = task_def.get("queue", "default")
589589

590590
from .queue import Queue
591-
queue_obj = Queue(queue)
591+
queue_obj = Queue(queue, add_to_known_queues=True)
592592

593593
if queue_obj.is_raw:
594594
raise Exception("Can't queue regular jobs on a raw queue")

mrq/queue.py

Lines changed: 29 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@ class Queue(object):
2121
# of Queue in the current process
2222
known_queues = {}
2323

24-
def __init__(self, queue_id, add_to_known_queues=True):
24+
def __init__(self, queue_id, add_to_known_queues=False):
2525

2626
if isinstance(queue_id, Queue):
2727
self.id = queue_id.id # TODO use __new__?
@@ -110,7 +110,7 @@ def redis_known_subqueues(self):
110110

111111
for key in Queue.known_queues:
112112
if key.startswith(self.id) and not key.endswith(delimiter):
113-
queues.append(Queue(key))
113+
queues.append(Queue(key, add_to_known_queues=True))
114114

115115
return queues
116116

@@ -235,11 +235,35 @@ def all_active(cls):
235235
return queues
236236

237237
@classmethod
238-
def all_known(cls, ):
238+
def all_known(cls):
239239
""" List all previously known queues """
240240

241-
# raw queues we know exist from the config + known queues in redis
242-
return set(context.get_current_config().get("raw_queues", {}).keys() + cls.redis_known_queues().keys())
241+
# queues we know exist from the config + known queues in redis
242+
return cls.all_known_from_config().union(set(cls.redis_known_queues().keys()))
243+
244+
@classmethod
245+
def all_known_from_config(cls):
246+
""" List all known queues from config (raw and regular). Caution: this does not account for the
247+
configuration of workers' queues (usually given via command line)
248+
"""
249+
250+
cfg = context.get_current_config()
251+
252+
queues_from_config = [
253+
t.get("queue")
254+
for t in (cfg.get("tasks") or {}).values()
255+
if t.get("queue")
256+
]
257+
258+
queues_from_config += (cfg.get("raw_queues") or {}).keys()
259+
260+
queues_from_config += [
261+
t.get("retry_queue")
262+
for t in (cfg.get("raw_queues") or {}).values()
263+
if t.get("retry_queue")
264+
]
265+
266+
return set(queues_from_config)
243267

244268
@classmethod
245269
def all(cls):

mrq/worker.py

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -73,7 +73,7 @@ def __init__(self):
7373
self.log_handler = LogHandler(quiet=self.config["quiet"])
7474
self.log = self.log_handler.get_logger(worker=self.id)
7575

76-
self.queues = [Queue(x) for x in self.config["queues"] if x]
76+
self.queues = [Queue(x, add_to_known_queues=True) for x in self.config["queues"] if x]
7777

7878
self.log.info(
7979
"Starting Gevent pool with %s worker greenlets (+ report, logs, adminhttp)" %
@@ -212,9 +212,9 @@ def greenlet_subqueues(self):
212212
try:
213213
for queue in self.config["queues"]:
214214
if queue.endswith(get_current_config().get("subqueues_delimiter")):
215-
queues += Queue(queue).redis_known_subqueues()
215+
queues += Queue(queue, add_to_known_queues=True).redis_known_subqueues()
216216
else:
217-
queues.append(Queue(queue))
217+
queues.append(Queue(queue, add_to_known_queues=True))
218218

219219
except Exception as e: # pylint: disable=broad-except
220220
self.log.error("When refreshing subqueues: %s", e)

0 commit comments

Comments
 (0)