Skip to content

Commit a3fc2ec

Browse files
committed
Factoring move with requeue (basetasks/utils.py)
1 parent 2c73b9b commit a3fc2ec

1 file changed

Lines changed: 2 additions & 25 deletions

File tree

mrq/basetasks/utils.py

Lines changed: 2 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -98,30 +98,7 @@ def perform_action(self, action, query, destination_queue):
9898
if list(query.keys()) == ["queue"]:
9999
Queue(query["queue"]).empty()
100100

101-
elif action == "move":
102-
cursor = self.collection.find(query, projection=["_id", "queue"])
103-
fetched_jobs = list(cursor)
104-
for jobs in group_iter(fetched_jobs, n=1000):
105-
jobs_by_queue = defaultdict(list)
106-
for job in jobs:
107-
jobs_by_queue[job["queue"]].append(job["_id"])
108-
stats["requeued"] += 1
109-
110-
for queue in jobs_by_queue:
111-
112-
updates = {
113-
"status": "queued",
114-
"datequeued": datetime.datetime.utcnow(),
115-
"dateupdated": datetime.datetime.utcnow(),
116-
"queue": destination_queue,
117-
"retry_count": 0
118-
}
119-
120-
self.collection.update({
121-
"_id": {"$in": jobs_by_queue[queue]}
122-
}, {"$set": updates}, multi=True)
123-
124-
elif action in ("requeue", "requeue_retry"):
101+
elif action in ("requeue", "requeue_retry", "move"):
125102

126103
# Requeue task by groups of maximum 1k items (if all in the same
127104
# queue)
@@ -150,7 +127,7 @@ def perform_action(self, action, query, destination_queue):
150127
if destination_queue is not None:
151128
updates["queue"] = destination_queue
152129

153-
if action == "requeue":
130+
if action in ("requeue", "move"):
154131
updates["retry_count"] = 0
155132

156133
self.collection.update({

0 commit comments

Comments
 (0)