Skip to content

Commit 0e7a91a

Browse files
authored
Merge pull request #135 from FlorianPerucki/master
add pause/resume queue methods
2 parents 61c9442 + 1ed40d1 commit 0e7a91a

6 files changed

Lines changed: 227 additions & 6 deletions

File tree

docs/command-line.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -56,6 +56,7 @@ You can pass additional configuration flags:
5656
- `--report_file`: Filepath of a json dump of the worker status. Disabled if none.
5757
- `--subqueues_refresh_interval`: Seconds between worker refreshes of the known subqueues.
5858
- `--subqueues_delimiter`: Delimiter between main queue and subqueue names.
59+
- `--paused_queues_refresh_interval`: Seconds between worker refreshes of the paused queues list.
5960
- `--admin_port`: Start an admin server on this port, if provided. Incompatible with --processes. Defaults to **0**
6061
- `--admin_ip`: IP for the admin server to listen on. Use "0.0.0.0" to allow access from outside. Defaults to **127.0.0.1**.
6162
- `--local_ip`: Overwrite the local IP, to be displayed in the dashboard.
@@ -103,4 +104,3 @@ $ mrq-run tasks.mylib.myfile.MyTask '{"param1": 1, "param2": True}'
103104
# Shorter syntax which casts all values as strings (equivalent to '{"param1": "1", "param2": "ok"}')
104105
$ mrq-run tasks.mylib.myfile.MyTask param1 1 param2 ok
105106
```
106-

mrq/config.py

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -336,11 +336,18 @@ def add_parser_args(parser, config_type):
336336

337337
parser.add_argument(
338338
'--subqueues_refresh_interval',
339-
default=60,
339+
default=10,
340340
action='store',
341341
type=float,
342342
help="Seconds between worker refreshes of the known subqueues")
343343

344+
parser.add_argument(
345+
'--paused_queues_refresh_interval',
346+
default=10,
347+
action='store',
348+
type=float,
349+
help="Seconds between worker refreshes of the paused queues list")
350+
344351
parser.add_argument(
345352
'--subqueues_delimiter',
346353
default='/',

mrq/queue.py

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ class Queue(object):
2525
# This is a mutable type so it is shared by all instances
2626
# of Queue in the current process
2727
known_queues = {}
28+
paused_queues = set()
2829

2930
def __init__(self, queue_id, add_to_known_queues=False):
3031

@@ -76,6 +77,11 @@ def redis_key_started(cls):
7677
""" Returns the global redis key used to store started job ids """
7778
return "%s:s:started" % context.get_current_config()["redis_prefix"]
7879

80+
@classmethod
81+
def redis_key_paused_queues(cls):
82+
""" Returns the redis key used to store this queue. """
83+
return "%s:s:paused" % (context.get_current_config()["redis_prefix"])
84+
7985
@classmethod
8086
def redis_key_known_queues(cls):
8187
""" Returns the global redis key used to store started job ids """
@@ -112,6 +118,11 @@ def redis_known_queues(cls):
112118
for value, score in context.connections.redis.zrange(cls.redis_key_known_queues(), 0, -1, withscores=True)
113119
}
114120

121+
@classmethod
122+
def redis_paused_queues(cls):
123+
""" Returns the set of currently paused queues """
124+
return context.connections.redis.smembers(cls.redis_key_paused_queues())
125+
115126
def redis_known_subqueues(self):
116127
""" Return the known subqueues of this queue as Queue objects. """
117128
delimiter = context.get_current_config()["subqueues_delimiter"]
@@ -147,6 +158,37 @@ def unserialize_job_ids(self, job_ids):
147158
else:
148159
return [x.encode('hex') for x in job_ids]
149160

161+
def _get_pausable_id(self):
162+
"""
163+
Get the queue id (either id or root_id) that should be used to pause/unpause the current queue
164+
TODO: handle subqueues with more than one level, e.g. "queue/subqueue/"
165+
"""
166+
queue = self.id
167+
delimiter = context.get_current_config().get("subqueues_delimiter")
168+
if delimiter is not None and self.id.endswith(delimiter):
169+
queue = self.root_id
170+
return queue
171+
172+
def pause(self):
173+
""" Adds this queue to the set of paused queues """
174+
context.connections.redis.sadd(Queue.redis_key_paused_queues(), self._get_pausable_id())
175+
176+
def is_paused(self):
177+
"""
178+
Returns wether the queue is paused or not.
179+
Warning: this does NOT ensure that the queue was effectively added to
180+
the set of paused queues. See the 'paused_queues_refresh_interval' option.
181+
"""
182+
root_is_paused = False
183+
if self.root_id != self.id:
184+
root_is_paused = context.connections.redis.sismember(Queue.redis_key_paused_queues(), self.root_id)
185+
186+
return root_is_paused or context.connections.redis.sismember(Queue.redis_key_paused_queues(), self.id)
187+
188+
def resume(self):
189+
""" Resumes a paused queue """
190+
context.connections.redis.srem(Queue.redis_key_paused_queues(), self._get_pausable_id())
191+
150192
def size(self):
151193
""" Returns the total number of jobs on the queue """
152194

@@ -472,6 +514,7 @@ def dequeue_jobs(self, max_jobs=1, job_class=None, worker=None):
472514
for j in job_data:
473515
j["status"] = "started"
474516
j["queue"] = retry_queue
517+
j["raw_queue"] = self.id
475518
if worker:
476519
j["worker"] = worker.id
477520

mrq/worker.py

Lines changed: 24 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -227,6 +227,14 @@ def greenlet_subqueues(self):
227227

228228
time.sleep(self.config["subqueues_refresh_interval"])
229229

230+
def greenlet_paused_queues(self):
231+
232+
while True:
233+
234+
# Update the process-local list of paused queues
235+
Queue.paused_queues = Queue.redis_paused_queues()
236+
time.sleep(self.config["paused_queues_refresh_interval"])
237+
230238
def get_memory(self):
231239
mmaps = self.process.get_memory_maps()
232240
mem = {
@@ -416,7 +424,13 @@ def work_init(self):
416424

417425
self.status = "started"
418426

419-
self.greenlets["subqueues"] = gevent.spawn(self.greenlet_subqueues)
427+
# An interval of 0 disables the refresh
428+
if self.config["subqueues_refresh_interval"] > 0:
429+
self.greenlets["subqueues"] = gevent.spawn(self.greenlet_subqueues)
430+
431+
# An interval of 0 disables the refresh
432+
if self.config["paused_queues_refresh_interval"] > 0:
433+
self.greenlets["paused_queues"] = gevent.spawn(self.greenlet_paused_queues)
420434

421435
self.greenlets["report"] = gevent.spawn(self.greenlet_report)
422436

@@ -470,9 +484,15 @@ def work_loop(self, max_jobs=None):
470484

471485
jobs = []
472486

473-
for queue_i in xrange(len(self.queues)):
487+
available_queues = [
488+
queue for queue in self.queues
489+
if queue.root_id not in Queue.paused_queues and
490+
queue.id not in Queue.paused_queues
491+
]
492+
493+
for queue_i in xrange(len(available_queues)):
474494

475-
queue = self.queues[(queue_i + queue_offset) % len(self.queues)]
495+
queue = available_queues[(queue_i + queue_offset) % len(available_queues)]
476496

477497
max_jobs_per_queue = free_pool_slots - len(jobs)
478498

@@ -481,7 +501,7 @@ def work_loop(self, max_jobs=None):
481501
break
482502

483503
if self.config["dequeue_strategy"] == "parallel":
484-
max_jobs_per_queue = max(1, int(max_jobs_per_queue / (len(self.queues) - queue_i)))
504+
max_jobs_per_queue = max(1, int(max_jobs_per_queue / (len(available_queues) - queue_i)))
485505

486506
jobs += queue.dequeue_jobs(
487507
max_jobs=max_jobs_per_queue,

tests/test_pause.py

Lines changed: 130 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,130 @@
1+
from mrq.job import Job
2+
import pytest
3+
from mrq.queue import Queue, send_task
4+
import time
5+
from mrq.context import set_current_config, get_config
6+
7+
8+
def test_pause_resume(worker):
9+
10+
worker.start(flags="--paused_queues_refresh_interval=0.1")
11+
12+
Queue("high").pause()
13+
14+
assert Queue("high").is_paused()
15+
16+
# wait for the paused_queues list to be refreshed
17+
time.sleep(2)
18+
19+
job_id1 = send_task(
20+
"tests.tasks.general.MongoInsert", {"a": 41},
21+
queue="high")
22+
23+
job_id2 = send_task(
24+
"tests.tasks.general.MongoInsert", {"a": 43},
25+
queue="low")
26+
27+
time.sleep(5)
28+
29+
job1 = Job(job_id1).fetch().data
30+
job2 = Job(job_id2).fetch().data
31+
32+
assert job1["status"] == "queued"
33+
34+
assert job2["status"] == "success"
35+
assert job2["result"] == {"a": 43}
36+
37+
assert worker.mongodb_jobs.tests_inserts.count() == 1
38+
39+
Queue("high").resume()
40+
41+
Job(job_id1).wait(poll_interval=0.01)
42+
43+
job1 = Job(job_id1).fetch().data
44+
45+
assert job1["status"] == "success"
46+
assert job1["result"] == {"a": 41}
47+
48+
assert worker.mongodb_jobs.tests_inserts.count() == 2
49+
50+
worker.stop()
51+
52+
53+
def test_pause_refresh_interval(worker):
54+
55+
""" Tests that a refresh interval of 0 disables the pause functionnality """
56+
57+
worker.start(flags="--paused_queues_refresh_interval=0")
58+
59+
Queue("high").pause()
60+
61+
assert Queue("high").is_paused()
62+
63+
# wait for the paused_queues list to be refreshed
64+
time.sleep(2)
65+
66+
job_id1 = send_task(
67+
"tests.tasks.general.MongoInsert", {"a": 41},
68+
queue="high")
69+
70+
time.sleep(5)
71+
72+
job1 = Job(job_id1).fetch().data
73+
74+
assert job1["status"] == "success"
75+
assert job1["result"] == {"a": 41}
76+
77+
worker.stop()
78+
79+
80+
def test_pause_subqueue(worker):
81+
82+
# set config in current context in order to have a subqueue delimiter
83+
set_current_config(get_config(config_type="worker"))
84+
85+
worker.start(queues="high high/", flags="--subqueues_refresh_interval=1 --paused_queues_refresh_interval=1")
86+
87+
Queue("high").pause()
88+
89+
assert Queue("high/").is_paused()
90+
91+
# wait for the paused_queues list to be refreshed
92+
time.sleep(2)
93+
94+
job_id1 = send_task(
95+
"tests.tasks.general.MongoInsert", {"a": 41},
96+
queue="high")
97+
98+
job_id2 = send_task(
99+
"tests.tasks.general.MongoInsert", {"a": 43},
100+
queue="high/subqueue")
101+
102+
# wait a bit to make sure the jobs status will still be queued
103+
time.sleep(5)
104+
105+
job1 = Job(job_id1).fetch().data
106+
job2 = Job(job_id2).fetch().data
107+
108+
assert job1["status"] == "queued"
109+
assert job2["status"] == "queued"
110+
111+
assert worker.mongodb_jobs.tests_inserts.count() == 0
112+
113+
Queue("high/").resume()
114+
115+
Job(job_id1).wait(poll_interval=0.01)
116+
117+
Job(job_id2).wait(poll_interval=0.01)
118+
119+
job1 = Job(job_id1).fetch().data
120+
job2 = Job(job_id2).fetch().data
121+
122+
assert job1["status"] == "success"
123+
assert job1["result"] == {"a": 41}
124+
125+
assert job2["status"] == "success"
126+
assert job2["result"] == {"a": 43}
127+
128+
assert worker.mongodb_jobs.tests_inserts.count() == 2
129+
130+
worker.stop()

tests/test_subqueues.py

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -53,3 +53,24 @@ def test_custom_delimiters(worker, delimiter):
5353
job_id = worker.send_task("tests.tasks.general.GetTime", {}, queue=subqueue, block=False)
5454
Job(job_id).wait(poll_interval=0.01)
5555
worker.stop()
56+
57+
58+
def test_refresh_interval(worker):
59+
60+
""" Tests that a refresh interval of 0 disables the subqueue detection """
61+
62+
worker.start(queues="test/", flags="--subqueues_refresh_interval=0")
63+
64+
time.sleep(2)
65+
66+
job_id1 = worker.send_task(
67+
"tests.tasks.general.GetTime", {"a": 41},
68+
queue="test/subqueue", block=False)
69+
70+
time.sleep(5)
71+
72+
job1 = Job(job_id1).fetch().data
73+
74+
assert job1["status"] == "queued"
75+
76+
worker.stop()

0 commit comments

Comments
 (0)