From 7c95f7761f027e0742b843940e6e80f2a4d23c1a Mon Sep 17 00:00:00 2001 From: Florian Perucki Date: Wed, 19 Oct 2016 09:06:00 +0200 Subject: [PATCH 1/6] add pause/resume queue methods --- mrq/queue.py | 22 ++++++++++++++++++++++ mrq/worker.py | 9 ++++++--- tests/test_pause.py | 43 +++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 71 insertions(+), 3 deletions(-) create mode 100644 tests/test_pause.py diff --git a/mrq/queue.py b/mrq/queue.py index 61a04197..1c0a614c 100644 --- a/mrq/queue.py +++ b/mrq/queue.py @@ -76,6 +76,11 @@ def redis_key_started(cls): """ Returns the global redis key used to store started job ids """ return "%s:s:started" % context.get_current_config()["redis_prefix"] + @classmethod + def redis_key_paused_queues(cls): + """ Returns the redis key used to store this queue. """ + return "%s:s:paused" % (context.get_current_config()["redis_prefix"]) + @classmethod def redis_key_known_queues(cls): """ Returns the global redis key used to store started job ids """ @@ -112,6 +117,11 @@ def redis_known_queues(cls): for value, score in context.connections.redis.zrange(cls.redis_key_known_queues(), 0, -1, withscores=True) } + @classmethod + def redis_paused_queues(cls): + """ Returns the set of currently paused queues """ + return context.connections.redis.smembers(cls.redis_key_paused_queues()) + def redis_known_subqueues(self): """ Return the known subqueues of this queue as Queue objects. """ delimiter = context.get_current_config()["subqueues_delimiter"] @@ -147,6 +157,18 @@ def unserialize_job_ids(self, job_ids): else: return [x.encode('hex') for x in job_ids] + def pause(self): + """ Adds this queue to the set of paused queues """ + context.connections.redis.sadd(Queue.redis_key_paused_queues(), self.id) + + def is_paused(self): + """ Returns wether the queue is paused or not """ + return context.connections.redis.sismember(Queue.redis_key_paused_queues(), self.id) + + def resume(self): + """ Resumes a paused queue """ + context.connections.redis.srem(Queue.redis_key_paused_queues(), self.id) + def size(self): """ Returns the total number of jobs on the queue """ diff --git a/mrq/worker.py b/mrq/worker.py index f959a2bd..4d4d9072 100644 --- a/mrq/worker.py +++ b/mrq/worker.py @@ -207,6 +207,7 @@ def greenlet_subqueues(self): # Update the process-local list of known queues Queue.known_queues = Queue.redis_known_queues() + Queue.paused_queues = Queue.redis_paused_queues() queues = [] try: @@ -469,10 +470,12 @@ def work_loop(self, max_jobs=None): gevent.sleep(0.01) jobs = [] + paused_queues = Queue.redis_paused_queues() + available_queues = [queue for queue in self.queues if queue.id not in paused_queues] - for queue_i in xrange(len(self.queues)): + for queue_i in xrange(len(available_queues)): - queue = self.queues[(queue_i + queue_offset) % len(self.queues)] + queue = available_queues[(queue_i + queue_offset) % len(available_queues)] max_jobs_per_queue = free_pool_slots - len(jobs) @@ -481,7 +484,7 @@ def work_loop(self, max_jobs=None): break if self.config["dequeue_strategy"] == "parallel": - max_jobs_per_queue = max(1, int(max_jobs_per_queue / (len(self.queues) - queue_i))) + max_jobs_per_queue = max(1, int(max_jobs_per_queue / (len(available_queues) - queue_i))) jobs += queue.dequeue_jobs( max_jobs=max_jobs_per_queue, diff --git a/tests/test_pause.py b/tests/test_pause.py new file mode 100644 index 00000000..ae4c4fbb --- /dev/null +++ b/tests/test_pause.py @@ -0,0 +1,43 @@ +from mrq.job import Job +from mrq.queue import Queue, send_task +import time + + +def test_pause_resume(worker): + + worker.start() + + Queue("high").pause() + + assert Queue("high").is_paused() + + job_id1 = send_task( + "tests.tasks.general.MongoInsert", {"a": 41}, + queue="high") + + job_id2 = send_task( + "tests.tasks.general.MongoInsert", {"a": 43}, + queue="low") + + time.sleep(5) + + job1 = Job(job_id1).fetch().data + job2 = Job(job_id2).fetch().data + + assert job1["status"] == "queued" + + assert job2["status"] == "success" + assert job2["result"] == {"a": 43} + + assert worker.mongodb_jobs.tests_inserts.count() == 1 + + Queue("high").resume() + + Job(job_id1).wait(poll_interval=0.01) + + job1 = Job(job_id1).fetch().data + + assert job1["status"] == "success" + assert job1["result"] == {"a": 41} + + assert worker.mongodb_jobs.tests_inserts.count() == 2 From 5bb0f67b243e39038062e60d7bf416b9a0f9d1dd Mon Sep 17 00:00:00 2001 From: Florian Perucki Date: Wed, 19 Oct 2016 14:31:57 +0200 Subject: [PATCH 2/6] add a greenlet to refresh the list of paused queues --- docs/command-line.md | 2 +- mrq/config.py | 7 +++++++ mrq/queue.py | 7 ++++++- mrq/worker.py | 16 +++++++++++++--- tests/test_pause.py | 34 +++++++++++++++++++++++++++++++++- 5 files changed, 60 insertions(+), 6 deletions(-) diff --git a/docs/command-line.md b/docs/command-line.md index f5669cec..786ad2e0 100644 --- a/docs/command-line.md +++ b/docs/command-line.md @@ -56,6 +56,7 @@ You can pass additional configuration flags: - `--report_file`: Filepath of a json dump of the worker status. Disabled if none. - `--subqueues_refresh_interval`: Seconds between worker refreshes of the known subqueues. - `--subqueues_delimiter`: Delimiter between main queue and subqueue names. + - `--paused_queues_refresh_interval`: Seconds between worker refreshes of the paused queues list. - `--admin_port`: Start an admin server on this port, if provided. Incompatible with --processes. Defaults to **0** - `--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**. - `--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}' # Shorter syntax which casts all values as strings (equivalent to '{"param1": "1", "param2": "ok"}') $ mrq-run tasks.mylib.myfile.MyTask param1 1 param2 ok ``` - diff --git a/mrq/config.py b/mrq/config.py index ad772fab..feb44368 100644 --- a/mrq/config.py +++ b/mrq/config.py @@ -341,6 +341,13 @@ def add_parser_args(parser, config_type): type=float, help="Seconds between worker refreshes of the known subqueues") + parser.add_argument( + '--paused_queues_refresh_interval', + default=60, + action='store', + type=float, + help="Seconds between worker refreshes of the paused queues list") + parser.add_argument( '--subqueues_delimiter', default='/', diff --git a/mrq/queue.py b/mrq/queue.py index 1c0a614c..f9d6d583 100644 --- a/mrq/queue.py +++ b/mrq/queue.py @@ -25,6 +25,7 @@ class Queue(object): # This is a mutable type so it is shared by all instances # of Queue in the current process known_queues = {} + paused_queues = set() def __init__(self, queue_id, add_to_known_queues=False): @@ -162,7 +163,11 @@ def pause(self): context.connections.redis.sadd(Queue.redis_key_paused_queues(), self.id) def is_paused(self): - """ Returns wether the queue is paused or not """ + """ + Returns wether the queue is paused or not. + Warning: this does NOT ensure that the queue was effectively added to + the list of paused queues. See the 'paused_queues_refresh_interval' option. + """ return context.connections.redis.sismember(Queue.redis_key_paused_queues(), self.id) def resume(self): diff --git a/mrq/worker.py b/mrq/worker.py index 4d4d9072..7e0b6a13 100644 --- a/mrq/worker.py +++ b/mrq/worker.py @@ -207,7 +207,6 @@ def greenlet_subqueues(self): # Update the process-local list of known queues Queue.known_queues = Queue.redis_known_queues() - Queue.paused_queues = Queue.redis_paused_queues() queues = [] try: @@ -228,6 +227,14 @@ def greenlet_subqueues(self): time.sleep(self.config["subqueues_refresh_interval"]) + def greenlet_paused_queues(self): + + while True: + + # Update the process-local list of known queues + Queue.paused_queues = Queue.redis_paused_queues() + time.sleep(self.config["paused_queues_refresh_interval"]) + def get_memory(self): mmaps = self.process.get_memory_maps() mem = { @@ -419,6 +426,10 @@ def work_init(self): self.greenlets["subqueues"] = gevent.spawn(self.greenlet_subqueues) + # An interval of 0 disables the refresh + if self.config["paused_queues_refresh_interval"] > 0: + self.greenlets["paused_queues"] = gevent.spawn(self.greenlet_paused_queues) + self.greenlets["report"] = gevent.spawn(self.greenlet_report) self.greenlets["logs"] = gevent.spawn(self.greenlet_logs) @@ -470,8 +481,7 @@ def work_loop(self, max_jobs=None): gevent.sleep(0.01) jobs = [] - paused_queues = Queue.redis_paused_queues() - available_queues = [queue for queue in self.queues if queue.id not in paused_queues] + available_queues = [queue for queue in self.queues if queue.id not in Queue.paused_queues] for queue_i in xrange(len(available_queues)): diff --git a/tests/test_pause.py b/tests/test_pause.py index ae4c4fbb..006f9b49 100644 --- a/tests/test_pause.py +++ b/tests/test_pause.py @@ -5,12 +5,15 @@ def test_pause_resume(worker): - worker.start() + worker.start(flags="--paused_queues_refresh_interval=0.1") Queue("high").pause() assert Queue("high").is_paused() + # wait for the paused_queues list to be refreshed + time.sleep(2) + job_id1 = send_task( "tests.tasks.general.MongoInsert", {"a": 41}, queue="high") @@ -41,3 +44,32 @@ def test_pause_resume(worker): assert job1["result"] == {"a": 41} assert worker.mongodb_jobs.tests_inserts.count() == 2 + + worker.stop() + + +def test_pause_refresh_interval(worker): + + """ Tests that a refresh interval of 0 disables the pause functionnality """ + + worker.start(flags="--paused_queues_refresh_interval=0") + + Queue("high").pause() + + assert Queue("high").is_paused() + + # wait for the paused_queues list to be refreshed + time.sleep(2) + + job_id1 = send_task( + "tests.tasks.general.MongoInsert", {"a": 41}, + queue="high") + + time.sleep(5) + + job1 = Job(job_id1).fetch().data + + assert job1["status"] == "success" + assert job1["result"] == {"a": 41} + + worker.stop() From d77a49e8c6bc3122b39cb005cf6cdca17c0ef279 Mon Sep 17 00:00:00 2001 From: Florian Perucki Date: Wed, 19 Oct 2016 14:40:46 +0200 Subject: [PATCH 3/6] a refresh interval of 0 should disable the subqueues detection --- mrq/worker.py | 4 +++- tests/test_subqueues.py | 21 +++++++++++++++++++++ 2 files changed, 24 insertions(+), 1 deletion(-) diff --git a/mrq/worker.py b/mrq/worker.py index 7e0b6a13..efee401f 100644 --- a/mrq/worker.py +++ b/mrq/worker.py @@ -424,7 +424,9 @@ def work_init(self): self.status = "started" - self.greenlets["subqueues"] = gevent.spawn(self.greenlet_subqueues) + # An interval of 0 disables the refresh + if self.config["subqueues_refresh_interval"] > 0: + self.greenlets["subqueues"] = gevent.spawn(self.greenlet_subqueues) # An interval of 0 disables the refresh if self.config["paused_queues_refresh_interval"] > 0: diff --git a/tests/test_subqueues.py b/tests/test_subqueues.py index 34b65f33..4c58561f 100644 --- a/tests/test_subqueues.py +++ b/tests/test_subqueues.py @@ -53,3 +53,24 @@ def test_custom_delimiters(worker, delimiter): job_id = worker.send_task("tests.tasks.general.GetTime", {}, queue=subqueue, block=False) Job(job_id).wait(poll_interval=0.01) worker.stop() + + +def test_refresh_interval(worker): + + """ Tests that a refresh interval of 0 disables the subqueue detection """ + + worker.start(queues="test/", flags="--subqueues_refresh_interval=0") + + time.sleep(2) + + job_id1 = worker.send_task( + "tests.tasks.general.GetTime", {"a": 41}, + queue="test/subqueue", block=False) + + time.sleep(5) + + job1 = Job(job_id1).fetch().data + + assert job1["status"] == "queued" + + worker.stop() From 3f23427e3f473506bfd5bcd828c3b7cf4bd91aed Mon Sep 17 00:00:00 2001 From: Florian Perucki Date: Thu, 20 Oct 2016 15:50:03 +0200 Subject: [PATCH 4/6] smaller refresh_interval default for subqueues and paused queues --- mrq/config.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/mrq/config.py b/mrq/config.py index feb44368..044ce7cf 100644 --- a/mrq/config.py +++ b/mrq/config.py @@ -336,14 +336,14 @@ def add_parser_args(parser, config_type): parser.add_argument( '--subqueues_refresh_interval', - default=60, + default=10, action='store', type=float, help="Seconds between worker refreshes of the known subqueues") parser.add_argument( '--paused_queues_refresh_interval', - default=60, + default=10, action='store', type=float, help="Seconds between worker refreshes of the paused queues list") From 1638a89791fe613edd9e3a2b420470aa36aee5ca Mon Sep 17 00:00:00 2001 From: Florian Perucki Date: Thu, 20 Oct 2016 18:54:36 +0200 Subject: [PATCH 5/6] add initial queue info in raw task jobs --- mrq/queue.py | 1 + 1 file changed, 1 insertion(+) diff --git a/mrq/queue.py b/mrq/queue.py index f9d6d583..0dd5ac1c 100644 --- a/mrq/queue.py +++ b/mrq/queue.py @@ -499,6 +499,7 @@ def dequeue_jobs(self, max_jobs=1, job_class=None, worker=None): for j in job_data: j["status"] = "started" j["queue"] = retry_queue + j["raw_queue"] = self.id if worker: j["worker"] = worker.id From 1ed40d1561c2f291145af41e41592e2f722d0034 Mon Sep 17 00:00:00 2001 From: Florian Perucki Date: Fri, 21 Oct 2016 14:24:15 +0200 Subject: [PATCH 6/6] pause subqueues and the root queue is paused --- mrq/queue.py | 23 +++++++++++++++---- mrq/worker.py | 9 ++++++-- tests/test_pause.py | 55 +++++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 81 insertions(+), 6 deletions(-) diff --git a/mrq/queue.py b/mrq/queue.py index 0dd5ac1c..9803fad5 100644 --- a/mrq/queue.py +++ b/mrq/queue.py @@ -158,21 +158,36 @@ def unserialize_job_ids(self, job_ids): else: return [x.encode('hex') for x in job_ids] + def _get_pausable_id(self): + """ + Get the queue id (either id or root_id) that should be used to pause/unpause the current queue + TODO: handle subqueues with more than one level, e.g. "queue/subqueue/" + """ + queue = self.id + delimiter = context.get_current_config().get("subqueues_delimiter") + if delimiter is not None and self.id.endswith(delimiter): + queue = self.root_id + return queue + def pause(self): """ Adds this queue to the set of paused queues """ - context.connections.redis.sadd(Queue.redis_key_paused_queues(), self.id) + context.connections.redis.sadd(Queue.redis_key_paused_queues(), self._get_pausable_id()) def is_paused(self): """ Returns wether the queue is paused or not. Warning: this does NOT ensure that the queue was effectively added to - the list of paused queues. See the 'paused_queues_refresh_interval' option. + the set of paused queues. See the 'paused_queues_refresh_interval' option. """ - return context.connections.redis.sismember(Queue.redis_key_paused_queues(), self.id) + root_is_paused = False + if self.root_id != self.id: + root_is_paused = context.connections.redis.sismember(Queue.redis_key_paused_queues(), self.root_id) + + return root_is_paused or context.connections.redis.sismember(Queue.redis_key_paused_queues(), self.id) def resume(self): """ Resumes a paused queue """ - context.connections.redis.srem(Queue.redis_key_paused_queues(), self.id) + context.connections.redis.srem(Queue.redis_key_paused_queues(), self._get_pausable_id()) def size(self): """ Returns the total number of jobs on the queue """ diff --git a/mrq/worker.py b/mrq/worker.py index efee401f..51bcaeae 100644 --- a/mrq/worker.py +++ b/mrq/worker.py @@ -231,7 +231,7 @@ def greenlet_paused_queues(self): while True: - # Update the process-local list of known queues + # Update the process-local list of paused queues Queue.paused_queues = Queue.redis_paused_queues() time.sleep(self.config["paused_queues_refresh_interval"]) @@ -483,7 +483,12 @@ def work_loop(self, max_jobs=None): gevent.sleep(0.01) jobs = [] - available_queues = [queue for queue in self.queues if queue.id not in Queue.paused_queues] + + available_queues = [ + queue for queue in self.queues + if queue.root_id not in Queue.paused_queues and + queue.id not in Queue.paused_queues + ] for queue_i in xrange(len(available_queues)): diff --git a/tests/test_pause.py b/tests/test_pause.py index 006f9b49..9ee0dd41 100644 --- a/tests/test_pause.py +++ b/tests/test_pause.py @@ -1,6 +1,8 @@ from mrq.job import Job +import pytest from mrq.queue import Queue, send_task import time +from mrq.context import set_current_config, get_config def test_pause_resume(worker): @@ -73,3 +75,56 @@ def test_pause_refresh_interval(worker): assert job1["result"] == {"a": 41} worker.stop() + + +def test_pause_subqueue(worker): + + # set config in current context in order to have a subqueue delimiter + set_current_config(get_config(config_type="worker")) + + worker.start(queues="high high/", flags="--subqueues_refresh_interval=1 --paused_queues_refresh_interval=1") + + Queue("high").pause() + + assert Queue("high/").is_paused() + + # wait for the paused_queues list to be refreshed + time.sleep(2) + + job_id1 = send_task( + "tests.tasks.general.MongoInsert", {"a": 41}, + queue="high") + + job_id2 = send_task( + "tests.tasks.general.MongoInsert", {"a": 43}, + queue="high/subqueue") + + # wait a bit to make sure the jobs status will still be queued + time.sleep(5) + + job1 = Job(job_id1).fetch().data + job2 = Job(job_id2).fetch().data + + assert job1["status"] == "queued" + assert job2["status"] == "queued" + + assert worker.mongodb_jobs.tests_inserts.count() == 0 + + Queue("high/").resume() + + Job(job_id1).wait(poll_interval=0.01) + + Job(job_id2).wait(poll_interval=0.01) + + job1 = Job(job_id1).fetch().data + job2 = Job(job_id2).fetch().data + + assert job1["status"] == "success" + assert job1["result"] == {"a": 41} + + assert job2["status"] == "success" + assert job2["result"] == {"a": 43} + + assert worker.mongodb_jobs.tests_inserts.count() == 2 + + worker.stop()