Skip to content

Commit 5bb0f67

Browse files
author
Florian Perucki
committed
add a greenlet to refresh the list of paused queues
1 parent 7c95f77 commit 5bb0f67

5 files changed

Lines changed: 60 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: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -341,6 +341,13 @@ def add_parser_args(parser, config_type):
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=60,
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: 6 additions & 1 deletion
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

@@ -162,7 +163,11 @@ def pause(self):
162163
context.connections.redis.sadd(Queue.redis_key_paused_queues(), self.id)
163164

164165
def is_paused(self):
165-
""" Returns wether the queue is paused or not """
166+
"""
167+
Returns wether the queue is paused or not.
168+
Warning: this does NOT ensure that the queue was effectively added to
169+
the list of paused queues. See the 'paused_queues_refresh_interval' option.
170+
"""
166171
return context.connections.redis.sismember(Queue.redis_key_paused_queues(), self.id)
167172

168173
def resume(self):

mrq/worker.py

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -207,7 +207,6 @@ def greenlet_subqueues(self):
207207

208208
# Update the process-local list of known queues
209209
Queue.known_queues = Queue.redis_known_queues()
210-
Queue.paused_queues = Queue.redis_paused_queues()
211210

212211
queues = []
213212
try:
@@ -228,6 +227,14 @@ def greenlet_subqueues(self):
228227

229228
time.sleep(self.config["subqueues_refresh_interval"])
230229

230+
def greenlet_paused_queues(self):
231+
232+
while True:
233+
234+
# Update the process-local list of known queues
235+
Queue.paused_queues = Queue.redis_paused_queues()
236+
time.sleep(self.config["paused_queues_refresh_interval"])
237+
231238
def get_memory(self):
232239
mmaps = self.process.get_memory_maps()
233240
mem = {
@@ -419,6 +426,10 @@ def work_init(self):
419426

420427
self.greenlets["subqueues"] = gevent.spawn(self.greenlet_subqueues)
421428

429+
# An interval of 0 disables the refresh
430+
if self.config["paused_queues_refresh_interval"] > 0:
431+
self.greenlets["paused_queues"] = gevent.spawn(self.greenlet_paused_queues)
432+
422433
self.greenlets["report"] = gevent.spawn(self.greenlet_report)
423434

424435
self.greenlets["logs"] = gevent.spawn(self.greenlet_logs)
@@ -470,8 +481,7 @@ def work_loop(self, max_jobs=None):
470481
gevent.sleep(0.01)
471482

472483
jobs = []
473-
paused_queues = Queue.redis_paused_queues()
474-
available_queues = [queue for queue in self.queues if queue.id not in paused_queues]
484+
available_queues = [queue for queue in self.queues if queue.id not in Queue.paused_queues]
475485

476486
for queue_i in xrange(len(available_queues)):
477487

tests/test_pause.py

Lines changed: 33 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,12 +5,15 @@
55

66
def test_pause_resume(worker):
77

8-
worker.start()
8+
worker.start(flags="--paused_queues_refresh_interval=0.1")
99

1010
Queue("high").pause()
1111

1212
assert Queue("high").is_paused()
1313

14+
# wait for the paused_queues list to be refreshed
15+
time.sleep(2)
16+
1417
job_id1 = send_task(
1518
"tests.tasks.general.MongoInsert", {"a": 41},
1619
queue="high")
@@ -41,3 +44,32 @@ def test_pause_resume(worker):
4144
assert job1["result"] == {"a": 41}
4245

4346
assert worker.mongodb_jobs.tests_inserts.count() == 2
47+
48+
worker.stop()
49+
50+
51+
def test_pause_refresh_interval(worker):
52+
53+
""" Tests that a refresh interval of 0 disables the pause functionnality """
54+
55+
worker.start(flags="--paused_queues_refresh_interval=0")
56+
57+
Queue("high").pause()
58+
59+
assert Queue("high").is_paused()
60+
61+
# wait for the paused_queues list to be refreshed
62+
time.sleep(2)
63+
64+
job_id1 = send_task(
65+
"tests.tasks.general.MongoInsert", {"a": 41},
66+
queue="high")
67+
68+
time.sleep(5)
69+
70+
job1 = Job(job_id1).fetch().data
71+
72+
assert job1["status"] == "success"
73+
assert job1["result"] == {"a": 41}
74+
75+
worker.stop()

0 commit comments

Comments
 (0)