Skip to content

Commit 724fe3d

Browse files
author
maelorn
committed
fix tests
1 parent a3f8389 commit 724fe3d

3 files changed

Lines changed: 41 additions & 20 deletions

File tree

mrq/basetasks/cleaning.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -57,7 +57,7 @@ def run(self, params):
5757
# iterate over them.
5858

5959
fields = {
60-
"_id": 1, "datestarted": 1, "queue": 1, "path": 1, "retry_count": 1, "worker": 1
60+
"_id": 1, "datestarted": 1, "queue": 1, "path": 1, "retry_count": 1, "worker": 1, "status": 1
6161
}
6262
for job_data in connections.mongodb_jobs.mrq_jobs.find(
6363
{"status": "started"}, projection=fields):

mrq/basetasks/utils.py

Lines changed: 27 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55
from mrq.queue import Queue
66
from bson import ObjectId
77
from mrq.context import connections, get_current_config, get_current_job
8+
from mrq.job import set_queues_size
89
from collections import defaultdict
910
from mrq.utils import group_iter
1011
import datetime
@@ -59,7 +60,7 @@ def build_query(self):
5960

6061
for key in params_dict:
6162
query["params.%s" % key] = params_dict[key]
62-
63+
6364
if current_job and "_id" not in query:
6465
query["_id"] = {"$lte": current_job.id}
6566

@@ -91,13 +92,24 @@ def perform_action(self, action, query, destination_queue):
9192
result_ttl = max([default_job_timeout] + tasks_ttls)
9293

9394
now = datetime.datetime.utcnow()
95+
96+
size_by_queues = defaultdict(int)
97+
if "queue" not in query:
98+
for job in self.collection.find(query, projection={"queue": 1}):
99+
size_by_queues[job["queue"]] += 1
100+
94101
ret = self.collection.update(query, {"$set": {
95102
"status": "cancel",
96103
"dateexpires": now + datetime.timedelta(seconds=result_ttl),
97104
"dateupdated": now
98105
}}, multi=True)
99106
stats["cancelled"] = ret["n"]
100107

108+
if "queue" in query:
109+
if isinstance(query["queue"], str):
110+
size_by_queues[query["queue"]] = ret["n"]
111+
set_queues_size(size_by_queues, action="decr")
112+
101113
# Special case when emptying just by queue name: empty it directly!
102114
# In this case we could also loose some jobs that were queued after
103115
# the MongoDB update. They will be "lost" and requeued later like the other case
@@ -123,22 +135,23 @@ def perform_action(self, action, query, destination_queue):
123135
jobs_by_queue[job["queue"]].append(job["_id"])
124136
stats["requeued"] += 1
125137

126-
for queue in jobs_by_queue:
138+
for queue in jobs_by_queue:
139+
updates = {
140+
"status": "queued",
141+
"datequeued": datetime.datetime.utcnow(),
142+
"dateupdated": datetime.datetime.utcnow()
143+
}
127144

128-
updates = {
129-
"status": "queued",
130-
"datequeued": datetime.datetime.utcnow(),
131-
"dateupdated": datetime.datetime.utcnow()
132-
}
145+
if destination_queue is not None:
146+
updates["queue"] = destination_queue
133147

134-
if destination_queue is not None:
135-
updates["queue"] = destination_queue
148+
if action == "requeue":
149+
updates["retry_count"] = 0
136150

137-
if action == "requeue":
138-
updates["retry_count"] = 0
151+
self.collection.update({
152+
"_id": {"$in": jobs_by_queue[queue]}
153+
}, {"$set": updates}, multi=True)
139154

140-
self.collection.update({
141-
"_id": {"$in": jobs_by_queue[queue]}
142-
}, {"$set": updates}, multi=True)
155+
set_queues_size({queue: len(jobs) for queue, jobs in jobs_by_queue.iteritems()})
143156

144157
return stats

mrq/job.py

Lines changed: 13 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -453,10 +453,12 @@ def _save_status(self, status, updates=None, exception=False, w=None, j=None):
453453
return
454454

455455
with context.connections.redis.pipeline(transaction=False) as pipe:
456-
if status != "started":
457-
if status == "queued":
456+
if status != "started" and "raw_queue" not in self.data:
457+
if status == "queued" and self.data["status"] != "started":
458+
print("INCR", self.data["queue"])
458459
pipe.incr("queuesize:%s" % self.data["queue"])
459460
else:
461+
print("DECR", self.data["queue"], self.data)
460462
pipe.decr("queuesize:%s" % self.data["queue"])
461463
pipe.expire("queuesize:%s" % self.data["queue"], context.get_current_config().get("queue_ttl"))
462464
pipe.execute()
@@ -634,8 +636,14 @@ def queue_job(main_task_path, params, **kwargs):
634636

635637
return queue_jobs(main_task_path, [params], **kwargs)[0]
636638

637-
def set_queue_size(queue, size):
638-
context.connections.redis.setex("queuesize:%s" % queue, context.get_current_config().get("queue_ttl"), size)
639+
def set_queues_size(size_by_queues, action="incr"):
640+
print("BLA", size_by_queues, action)
641+
if len(size_by_queues) > 0:
642+
with context.connections.redis.pipeline(transaction=False) as pipe:
643+
for queue in size_by_queues:
644+
getattr(pipe, action)("queuesize:%s" % queue, amount=size_by_queues[queue])
645+
pipe.expire("queuesize:%s" % queue, context.get_current_config().get("queue_ttl"))
646+
pipe.execute()
639647

640648
def queue_jobs(main_task_path, params_list, queue=None, batch_size=1000):
641649
""" Queue multiple jobs on a regular queue """
@@ -669,6 +677,6 @@ def queue_jobs(main_task_path, params_list, queue=None, batch_size=1000):
669677
all_ids += job_ids
670678

671679
queue_obj.notify(len(all_ids))
672-
set_queue_size(queue, len(all_ids))
680+
set_queues_size({queue: len(all_ids)})
673681

674682
return all_ids

0 commit comments

Comments
 (0)