Skip to content

Commit 4934e07

Browse files
author
maelorn
committed
refactor queuesize update
1 parent 724fe3d commit 4934e07

5 files changed

Lines changed: 39 additions & 26 deletions

File tree

mrq/basetasks/cleaning.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@ class RequeueRetryJobs(Task):
2727
max_concurrency = 1
2828

2929
def run(self, params):
30+
print("IN")
3031
return run_task("mrq.basetasks.utils.JobAction", {
3132
"status": "retry",
3233
"dateretry": {"$lte": datetime.datetime.utcnow()},

mrq/basetasks/utils.py

Lines changed: 16 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -135,22 +135,22 @@ def perform_action(self, action, query, destination_queue):
135135
jobs_by_queue[job["queue"]].append(job["_id"])
136136
stats["requeued"] += 1
137137

138-
for queue in jobs_by_queue:
139-
updates = {
140-
"status": "queued",
141-
"datequeued": datetime.datetime.utcnow(),
142-
"dateupdated": datetime.datetime.utcnow()
143-
}
144-
145-
if destination_queue is not None:
146-
updates["queue"] = destination_queue
147-
148-
if action == "requeue":
149-
updates["retry_count"] = 0
150-
151-
self.collection.update({
152-
"_id": {"$in": jobs_by_queue[queue]}
153-
}, {"$set": updates}, multi=True)
138+
for queue in jobs_by_queue:
139+
updates = {
140+
"status": "queued",
141+
"datequeued": datetime.datetime.utcnow(),
142+
"dateupdated": datetime.datetime.utcnow()
143+
}
144+
145+
if destination_queue is not None:
146+
updates["queue"] = destination_queue
147+
148+
if action == "requeue":
149+
updates["retry_count"] = 0
150+
151+
self.collection.update({
152+
"_id": {"$in": jobs_by_queue[queue]}
153+
}, {"$set": updates}, multi=True)
154154

155155
set_queues_size({queue: len(jobs) for queue, jobs in jobs_by_queue.iteritems()})
156156

mrq/job.py

Lines changed: 19 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -453,14 +453,26 @@ 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" and "raw_queue" not in self.data:
457-
if status == "queued" and self.data["status"] != "started":
458-
print("INCR", self.data["queue"])
459-
pipe.incr("queuesize:%s" % self.data["queue"])
460-
else:
461-
print("DECR", self.data["queue"], self.data)
456+
queue = (updates or {}).get("queue") or self.data["queue"]
457+
if status != "started":
458+
# Queue change
459+
if queue != self.data["queue"]:
462460
pipe.decr("queuesize:%s" % self.data["queue"])
463-
pipe.expire("queuesize:%s" % self.data["queue"], context.get_current_config().get("queue_ttl"))
461+
if status == "queued":
462+
pipe.incr("queuesize:%s" % queue)
463+
464+
# Regular queues
465+
elif status == "queued" and self.data.get("status") != "started":
466+
pipe.incr("queuesize:%s" % queue)
467+
468+
elif status != "queued" and not self.data.get("raw_queue"):
469+
pipe.decr("queuesize:%s" % queue)
470+
471+
# Raw queues retries
472+
elif (updates or {}).get("retry_count", 0) > self.data.get("retry_count", 0):
473+
pipe.incr("queuesize:%s" % queue)
474+
475+
pipe.expire("queuesize:%s" % queue, context.get_current_config().get("queue_ttl"))
464476
pipe.execute()
465477

466478
now = datetime.datetime.utcnow()
@@ -637,7 +649,6 @@ def queue_job(main_task_path, params, **kwargs):
637649
return queue_jobs(main_task_path, [params], **kwargs)[0]
638650

639651
def set_queues_size(size_by_queues, action="incr"):
640-
print("BLA", size_by_queues, action)
641652
if len(size_by_queues) > 0:
642653
with context.connections.redis.pipeline(transaction=False) as pipe:
643654
for queue in size_by_queues:

mrq/queue_raw.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -227,6 +227,7 @@ def dequeue_jobs(self, max_jobs=1, job_class=None, worker=None):
227227

228228
job_data = [job_factory(p) for p in params]
229229
for j in job_data:
230+
print("J")
230231
j["status"] = "started"
231232
j["queue"] = retry_queue
232233
j["datequeued"] = datetime.now()

tests/test_interrupts.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -192,7 +192,7 @@ def test_interrupt_worker_sigkill(worker, p_flags):
192192
"datestarted": datetime.datetime.utcnow() - datetime.timedelta(seconds=300)
193193
}})
194194

195-
assert Queue("default").size() == 0
195+
assert Queue("default").size() == 1
196196

197197
worker.start(queues="cleaning", deps=False, flush=False,
198198
flags=" --config tests/fixtures/config-shorttimeout.py")
@@ -203,7 +203,7 @@ def test_interrupt_worker_sigkill(worker, p_flags):
203203
assert res["requeued"] == 0
204204
assert res["started"] == 2 # current job should count too
205205

206-
assert Queue("default").size() == 0
206+
assert Queue("default").size() == 1
207207

208208
job = Job(job_id).fetch().data
209209
assert job["status"] == "started"

0 commit comments

Comments
 (0)