Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 22 additions & 0 deletions mrq/queue.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 """
Expand Down Expand Up @@ -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"]
Expand Down Expand Up @@ -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 """

Expand Down
9 changes: 6 additions & 3 deletions mrq/worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -469,10 +470,12 @@ def work_loop(self, max_jobs=None):
gevent.sleep(0.01)

jobs = []
paused_queues = Queue.redis_paused_queues()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

à passer dans un greenlet avec interval paramétrable, comme les subqueues

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

  • si intervale à zéro, couper le greenlet (et faire pareil dans subqueues?)

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)

Expand All @@ -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,
Expand Down
43 changes: 43 additions & 0 deletions tests/test_pause.py
Original file line number Diff line number Diff line change
@@ -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