Skip to content

Commit df43ea5

Browse files
committed
Fix tests & restore queued order
1 parent 9c1550e commit df43ea5

9 files changed

Lines changed: 85 additions & 26 deletions

File tree

mrq/basetasks/utils.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -120,6 +120,7 @@ def perform_action(self, action, query, destination_queue):
120120

121121
updates = {
122122
"status": "queued",
123+
"datequeued": datetime.datetime.utcnow(),
123124
"dateupdated": datetime.datetime.utcnow()
124125
}
125126

mrq/job.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -266,6 +266,7 @@ def requeue(self, queue=None, retry_count=0):
266266

267267
self._save_status("queued", updates={
268268
"queue": queue,
269+
"datequeued": datetime.datetime.utcnow(),
269270
"retry_count": retry_count
270271
})
271272

@@ -632,6 +633,7 @@ def queue_jobs(main_task_path, params_list, queue=None, batch_size=1000):
632633
"path": main_task_path,
633634
"params": params,
634635
"queue": queue,
636+
"datequeued": datetime.datetime.utcnow(),
635637
"status": "queued"
636638
} for params in params_group], w=1, return_jobs=False)
637639

mrq/monkey.py

Lines changed: 0 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -389,10 +389,6 @@ class mrq_patched_pymongo_cursor(Cursor):
389389
# Some dark magic is needed here to cope with python's name mangling for private variables.
390390
def _Cursor__send_message(self, *args, **kwargs):
391391

392-
if len(args) and args[0].name == "find":
393-
return Cursor._Cursor__send_message(self, *args, **kwargs)
394-
395-
# print self.__dict__
396392
job = get_current_job()
397393

398394
if job:

mrq/queue_raw.py

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -142,7 +142,6 @@ def dequeue_jobs(self, max_jobs=1, job_class=None, worker=None):
142142
retry_queue = self.get_retry_queue()
143143

144144
params = []
145-
jobs = []
146145

147146
# ZSET with times
148147
if self.is_timed:
@@ -187,10 +186,11 @@ def dequeue_jobs(self, max_jobs=1, job_class=None, worker=None):
187186
params = redis_group_command("lpop", max_jobs, self.redis_key)
188187

189188
if len(params) == 0:
190-
return []
189+
return
191190

192191
if worker:
193192
worker.status = "spawn"
193+
worker.idle_event.clear()
194194

195195
job_data = [job_factory(p) for p in params]
196196
for j in job_data:
@@ -200,9 +200,8 @@ def dequeue_jobs(self, max_jobs=1, job_class=None, worker=None):
200200
if worker:
201201
j["worker"] = worker.id
202202

203-
jobs += job_class.insert(job_data, statuses_no_storage=statuses_no_storage)
204-
205-
return jobs
203+
for job in job_class.insert(job_data, statuses_no_storage=statuses_no_storage):
204+
yield job
206205

207206
def get_sorted_graph(
208207
self,

mrq/queue_regular.py

Lines changed: 36 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@ def size(self):
1818
def list_job_ids(self, skip=0, limit=20):
1919
""" Returns a list of job ids on a queue """
2020

21-
return [str(x) for x in self.collection.find(
21+
return [str(x["_id"]) for x in self.collection.find(
2222
{"status": "queued"},
2323
sort=[("_id", -1 if self.is_reverse else 1)],
2424
projection={"_id": 1})
@@ -31,25 +31,49 @@ def dequeue_jobs(self, max_jobs=1, job_class=None, worker=None):
3131
from .job import Job
3232
job_class = Job
3333

34-
if worker:
35-
worker.status = "spawn"
36-
worker.idle_event.clear()
37-
3834
count = 0
3935

40-
for _ in range(max_jobs):
36+
job_ids = None
4137

42-
job_data = self.collection.find_one_and_update(
38+
# TODO: remove _id sort after full migration to datequeued
39+
sort_order = [("datequeued", -1 if self.is_reverse else 1), ("_id", -1 if self.is_reverse else 1)]
40+
41+
# MongoDB optimization: with many jobs it's faster to fetch the IDs first and do the atomic update second
42+
# Some jobs may have been stolen by another worker in the meantime but it's a balance (should we over-fetch?)
43+
if max_jobs > 5:
44+
job_ids = [x["_id"] for x in self.collection.find(
4345
{
4446
"status": "queued",
4547
"queue": self.id
4648
},
49+
limit=max_jobs,
50+
sort=sort_order,
51+
projection={"_id": 1}
52+
)]
53+
54+
if len(job_ids) == 0:
55+
return
56+
57+
for i in range(max_jobs if job_ids is None else len(job_ids)):
58+
59+
query = {
60+
"status": "queued",
61+
"queue": self.id
62+
}
63+
if job_ids is not None:
64+
query = {
65+
"status": "queued",
66+
"_id": job_ids[i]
67+
}
68+
69+
job_data = self.collection.find_one_and_update(
70+
query,
4771
{"$set": {
4872
"status": "started",
4973
"datestarted": datetime.datetime.utcnow(),
5074
"worker": worker.id if worker else None
5175
}},
52-
sort=[("_id", -1 if self.is_reverse else 1)],
76+
sort=sort_order,
5377
return_document=ReturnDocument.AFTER,
5478
projection={
5579
"_id": 1,
@@ -64,6 +88,10 @@ def dequeue_jobs(self, max_jobs=1, job_class=None, worker=None):
6488
if not job_data:
6589
break
6690

91+
if worker:
92+
worker.status = "spawn"
93+
worker.idle_event.clear()
94+
6795
count += 1
6896
context.metric("queues.%s.dequeued" % job_data["queue"], 1)
6997

mrq/worker.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -144,6 +144,8 @@ def ensure_indexes(self):
144144
[("dateexpires", 1)], sparse=True, background=False, expireAfterSeconds=0)
145145
self.mongodb_jobs.mrq_jobs.ensure_index(
146146
[("dateretry", 1)], sparse=True, background=False)
147+
self.mongodb_jobs.mrq_jobs.ensure_index(
148+
[("datequeued", 1)], sparse=True, background=True)
147149

148150
self.mongodb_jobs.mrq_scheduled_jobs.ensure_index(
149151
[("hash", 1)], unique=True, background=False, drop_dups=True)

tests/test_general.py

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -112,6 +112,26 @@ def test_general_simple_task_multiple(worker):
112112
[["dateupdated", 1]])] == [42, 42, 41]
113113

114114

115+
def test_general_requeue_order(worker):
116+
from mrq.job import Job
117+
118+
jobids = worker.send_tasks("tests.tasks.general.Add", [
119+
{"a": 41, "b": 1, "sleep": 4},
120+
{"a": 42, "b": 1, "sleep": 1},
121+
{"a": 43, "b": 1, "sleep": 1}
122+
], block=False)
123+
124+
time.sleep(2)
125+
126+
# We should be executing job1 now. Let's requeue job2, making it go to the end of the queue.
127+
Job(jobids[1]).requeue()
128+
129+
worker.get_wait_for_idle()
130+
131+
assert [x["result"] for x in worker.mongodb_jobs.mrq_jobs.find().sort(
132+
[["dateupdated", 1]])] == [42, 44, 43]
133+
134+
115135
def test_general_simple_task_reverse(worker):
116136

117137
worker.start(queues="default_reverse xtest test_timed_set", flags="--config tests/fixtures/config-raw1.py")

tests/test_io_hooks.py

Lines changed: 18 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -119,7 +119,7 @@ def test_io_hooks_mongodb(worker):
119119
assert job_events[4]["method"] == "cursor"
120120
assert job_events[5]["method"] == "cursor"
121121

122-
# Then getmore query
122+
# Then getmore query (can't understand why there are 2 more of those)
123123
assert job_events[6]["hook"] == "mongodb_pre"
124124
assert job_events[7]["hook"] == "mongodb_post"
125125

@@ -129,15 +129,25 @@ def test_io_hooks_mongodb(worker):
129129
assert job_events[6]["method"] == "cursor"
130130
assert job_events[7]["method"] == "cursor"
131131

132-
# Result MongoDB update
133-
132+
# Then getmore query
134133
assert job_events[8]["hook"] == "mongodb_pre"
135134
assert job_events[9]["hook"] == "mongodb_post"
136135

137-
assert job_events[8]["method"] == "update"
138-
assert job_events[9]["method"] == "update"
136+
assert job_events[8]["collection"] == "mrq.tests_inserts"
137+
assert job_events[9]["collection"] == "mrq.tests_inserts"
138+
139+
assert job_events[8]["method"] == "cursor"
140+
assert job_events[9]["method"] == "cursor"
141+
142+
# Result MongoDB update
143+
144+
assert job_events[10]["hook"] == "mongodb_pre"
145+
assert job_events[11]["hook"] == "mongodb_post"
146+
147+
assert job_events[10]["method"] == "update"
148+
assert job_events[11]["method"] == "update"
139149

140-
assert job_events[8]["collection"] == "mrq.mrq_jobs"
141-
assert job_events[9]["collection"] == "mrq.mrq_jobs"
150+
assert job_events[10]["collection"] == "mrq.mrq_jobs"
151+
assert job_events[11]["collection"] == "mrq.mrq_jobs"
142152

143-
assert len(job_events) == 5 * 2
153+
assert len(job_events) == 6 * 2

tests/test_jobinspect.py

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -102,9 +102,10 @@ def test_current_job_trace_io(worker, p_testtype, p_testparams, p_type, p_data,
102102
admin_worker = {}
103103
if len(admin_worker.get("jobs", [])) > 0:
104104
io = admin_worker["jobs"][0].get("io")
105+
105106
# Don't take MRQ's IOs as regular IO
106107
if io:
107-
if io["type"] == "mongodb" and io["data"]["collection"] in ["mrq.mrq_jobs", "mrq.mrq_logs"]:
108+
if io["type"].startswith("mongodb") and io["data"]["collection"] in ["mrq.mrq_jobs", "mrq.mrq_logs"]:
108109
io = False
109110
else:
110111
break

0 commit comments

Comments
 (0)