Skip to content

Commit 3b78c0e

Browse files
author
ismael
committed
Merge branch '0.9.x' into task_expiry
* 0.9.x: more readable condition fix cleaning logic make sure there's no job before cleaning a known queue Store subpool tracebacks + test fixes Fix Python 3 incompatibility (#188)
2 parents 5ee1335 + bfcc5bc commit 3b78c0e

10 files changed

Lines changed: 158 additions & 129 deletions

File tree

mrq/basetasks/cleaning.py

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -129,16 +129,19 @@ def run(self, params):
129129

130130
# Only clean queues older than N days
131131
time_threshold = time.time() - max_age
132-
for queue, time_last_used in known_queues.iteritems():
132+
for queue, time_last_used in known_queues.items():
133133
if queue in queues_from_config:
134134
continue
135135
if time_last_used < time_threshold:
136136
q = Queue(queue, add_to_known_queues=False)
137+
# size() returns the number of queued jobs
137138
size = q.size()
139+
140+
has_job = None
138141
if check_mongo:
139142
has_job = connections.mongodb_jobs.mrq_jobs.find_one({"queue": queue})
140143

141-
if size == 0 or has_job is None:
144+
if size == 0 and has_job is None:
142145
removed_queues.append(queue)
143146
print("Removing empty queue '%s' from known queues ..." % queue)
144147
if not pretend:

mrq/context.py

Lines changed: 2 additions & 116 deletions
Original file line numberDiff line numberDiff line change
@@ -10,8 +10,9 @@
1010
import pymongo
1111
import traceback
1212
from .utils import LazyObject, load_class_by_path
13-
from itertools import count as itertools_count
1413
from .config import get_config
14+
from .subpool import subpool_map, subpool_imap
15+
1516

1617
# This should be MRQ's only Python object shared by all the jobs in the same process
1718
_GLOBAL_CONTEXT = {
@@ -202,121 +203,6 @@ def trace(*args):
202203
greenlet.settrace(trace) # pylint: disable=no-member
203204

204205

205-
def subpool_map(pool_size, func, iterable):
206-
""" Starts a Gevent pool and run a map. Takes care of setting current_job and cleaning up. """
207-
208-
if not pool_size:
209-
return [func(*args) for args in iterable]
210-
211-
counter = itertools_count()
212-
213-
current_job = get_current_job()
214-
215-
def inner_func(*args):
216-
""" As each call to 'func' will be done in a random greenlet of the subpool, we need to
217-
register their IDs with set_current_job() to make get_current_job() calls work properly
218-
inside 'func'.
219-
"""
220-
next(counter)
221-
if current_job:
222-
set_current_job(current_job)
223-
224-
try:
225-
ret = func(*args)
226-
except Exception as exc:
227-
trace = traceback.format_exc()
228-
log.error("Error in subpool: %s \n%s" % (exc, trace))
229-
raise
230-
231-
if current_job:
232-
set_current_job(None)
233-
return ret
234-
235-
def inner_iterable():
236-
""" This will be called inside the pool's main greenlet, which ID also needs to be registered """
237-
if current_job:
238-
set_current_job(current_job)
239-
240-
for x in iterable:
241-
yield x
242-
243-
if current_job:
244-
set_current_job(None)
245-
246-
start_time = time.time()
247-
pool = gevent.pool.Pool(size=pool_size)
248-
ret = pool.map(inner_func, inner_iterable())
249-
pool.join(raise_error=True)
250-
total_time = time.time() - start_time
251-
252-
log.debug("SubPool ran %s greenlets in %0.6fs" % (counter, total_time))
253-
254-
return ret
255-
256-
257-
def subpool_imap(pool_size, func, iterable, flatten=False, unordered=False, buffer_size=None):
258-
""" Generator version of subpool_map. Should be used with unordered=True for optimal performance """
259-
260-
if not pool_size:
261-
for args in iterable:
262-
yield func(*args)
263-
264-
counter = itertools_count()
265-
266-
current_job = get_current_job()
267-
268-
def inner_func(*args):
269-
""" As each call to 'func' will be done in a random greenlet of the subpool, we need to
270-
register their IDs with set_current_job() to make get_current_job() calls work properly
271-
inside 'func'.
272-
"""
273-
next(counter)
274-
if current_job:
275-
set_current_job(current_job)
276-
277-
try:
278-
ret = func(*args)
279-
except Exception as exc:
280-
trace = traceback.format_exc()
281-
log.error("Error in subpool: %s \n%s" % (exc, trace))
282-
raise
283-
284-
if current_job:
285-
set_current_job(None)
286-
return ret
287-
288-
def inner_iterable():
289-
""" This will be called inside the pool's main greenlet, which ID also needs to be registered """
290-
if current_job:
291-
set_current_job(current_job)
292-
293-
for x in iterable:
294-
yield x
295-
296-
if current_job:
297-
set_current_job(None)
298-
299-
start_time = time.time()
300-
pool = gevent.pool.Pool(size=pool_size)
301-
302-
if unordered:
303-
iterator = pool.imap_unordered(inner_func, inner_iterable(), maxsize=buffer_size or pool_size)
304-
else:
305-
iterator = pool.imap(inner_func, inner_iterable())
306-
307-
for x in iterator:
308-
if flatten:
309-
for y in x:
310-
yield y
311-
else:
312-
yield x
313-
314-
pool.join(raise_error=True)
315-
total_time = time.time() - start_time
316-
317-
log.debug("SubPool ran %s greenlets in %0.6fs" % (counter, total_time))
318-
319-
320206
def run_task(path, params):
321207
""" Runs a task code synchronously """
322208
task_class = load_class_by_path(path)

mrq/job.py

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -473,8 +473,10 @@ def _save_status(self, status, updates=None, exception=False, w=None, j=None):
473473
if exception:
474474
trace = traceback.format_exc()
475475
context.log.error(trace)
476+
exc, value = sys.exc_info()[0:2]
477+
if hasattr(value, "subpool_traceback"):
478+
trace = "Exception first caught in a subpool. Traceback:\n%s\n%s" % (value.subpool_traceback, trace)
476479
db_updates["traceback"] = trace
477-
exc = sys.exc_info()[0]
478480
db_updates["exceptiontype"] = exc.__name__
479481

480482
self._save_traceback_history(status, trace, exc)

mrq/subpool.py

Lines changed: 123 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,123 @@
1+
from itertools import count as itertools_count
2+
import traceback
3+
import time
4+
import gevent
5+
6+
7+
def subpool_map(pool_size, func, iterable):
8+
""" Starts a Gevent pool and run a map. Takes care of setting current_job and cleaning up. """
9+
10+
from .context import get_current_job, set_current_job, log
11+
12+
if not pool_size:
13+
return [func(*args) for args in iterable]
14+
15+
counter = itertools_count()
16+
17+
current_job = get_current_job()
18+
19+
def inner_func(*args):
20+
""" As each call to 'func' will be done in a random greenlet of the subpool, we need to
21+
register their IDs with set_current_job() to make get_current_job() calls work properly
22+
inside 'func'.
23+
"""
24+
next(counter)
25+
if current_job:
26+
set_current_job(current_job)
27+
28+
try:
29+
ret = func(*args)
30+
except Exception as exc:
31+
trace = traceback.format_exc()
32+
exc.subpool_traceback = trace
33+
raise
34+
35+
if current_job:
36+
set_current_job(None)
37+
return ret
38+
39+
def inner_iterable():
40+
""" This will be called inside the pool's main greenlet, which ID also needs to be registered """
41+
if current_job:
42+
set_current_job(current_job)
43+
44+
for x in iterable:
45+
yield x
46+
47+
if current_job:
48+
set_current_job(None)
49+
50+
start_time = time.time()
51+
pool = gevent.pool.Pool(size=pool_size)
52+
ret = pool.map(inner_func, inner_iterable())
53+
pool.join(raise_error=True)
54+
total_time = time.time() - start_time
55+
56+
log.debug("SubPool ran %s greenlets in %0.6fs" % (counter, total_time))
57+
58+
return ret
59+
60+
61+
def subpool_imap(pool_size, func, iterable, flatten=False, unordered=False, buffer_size=None):
62+
""" Generator version of subpool_map. Should be used with unordered=True for optimal performance """
63+
64+
from .context import get_current_job, set_current_job, log
65+
66+
if not pool_size:
67+
for args in iterable:
68+
yield func(*args)
69+
70+
counter = itertools_count()
71+
72+
current_job = get_current_job()
73+
74+
def inner_func(*args):
75+
""" As each call to 'func' will be done in a random greenlet of the subpool, we need to
76+
register their IDs with set_current_job() to make get_current_job() calls work properly
77+
inside 'func'.
78+
"""
79+
next(counter)
80+
if current_job:
81+
set_current_job(current_job)
82+
83+
try:
84+
ret = func(*args)
85+
except Exception as exc:
86+
trace = traceback.format_exc()
87+
exc.subpool_traceback = trace
88+
raise
89+
90+
if current_job:
91+
set_current_job(None)
92+
return ret
93+
94+
def inner_iterable():
95+
""" This will be called inside the pool's main greenlet, which ID also needs to be registered """
96+
if current_job:
97+
set_current_job(current_job)
98+
99+
for x in iterable:
100+
yield x
101+
102+
if current_job:
103+
set_current_job(None)
104+
105+
start_time = time.time()
106+
pool = gevent.pool.Pool(size=pool_size)
107+
108+
if unordered:
109+
iterator = pool.imap_unordered(inner_func, inner_iterable(), maxsize=buffer_size or pool_size)
110+
else:
111+
iterator = pool.imap(inner_func, inner_iterable())
112+
113+
for x in iterator:
114+
if flatten:
115+
for y in x:
116+
yield y
117+
else:
118+
yield x
119+
120+
pool.join(raise_error=True)
121+
total_time = time.time() - start_time
122+
123+
log.debug("SubPool ran %s greenlets in %0.6fs" % (counter, total_time))

tests/conftest.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -261,7 +261,7 @@ def get_report(self, with_memory=False):
261261
f.close()
262262
return data
263263

264-
def wait_for_idle(self):
264+
def wait_for_idle(self, timeout=None):
265265

266266
if "--processes" in self.cmdline:
267267
print("Warning: wait_for_idle() doesn't support multiprocess workers yet")
@@ -275,7 +275,7 @@ def wait_for_idle(self):
275275
return True
276276

277277
try:
278-
wait_for_net_service("127.0.0.1", self.admin_port, poll_interval=0.01)
278+
wait_for_net_service("127.0.0.1", self.admin_port, poll_interval=0.01, timeout=timeout)
279279
f = urllib.request.urlopen("http://127.0.0.1:%s/wait_for_idle" % self.admin_port)
280280
data = f.read().decode('utf-8')
281281
assert data == "idle"

tests/tasks/general.py

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@
33
from builtins import range
44
from mrq.task import Task
55
from mrq.context import (log, retry_current_job, connections, get_current_config, get_current_job,
6-
subpool_map, abort_current_job, set_current_job_progress)
6+
subpool_map, subpool_imap, abort_current_job, set_current_job_progress)
77
from mrq.job import queue_job
88
import urllib.request, urllib.error, urllib.parse
99
import json
@@ -244,7 +244,7 @@ def inner(self, x):
244244
return True
245245

246246
if x == "exception":
247-
raise Exception(x)
247+
raise Exception(x) # __INNER_EXCEPTION_LINE__
248248

249249
time.sleep(x)
250250

@@ -253,7 +253,10 @@ def inner(self, x):
253253
def run(self, params):
254254
self.job = get_current_job()
255255

256-
return subpool_map(params["pool_size"], self.inner, params["inner_params"])
256+
if params.get("imap"):
257+
return subpool_map(params["pool_size"], self.inner, params["inner_params"])
258+
else:
259+
return list(subpool_imap(params["pool_size"], self.inner, params["inner_params"]))
257260

258261

259262
class GetMetrics(Task):

tests/test_agent.py

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -11,8 +11,9 @@ def scenario(profiles, agents):
1111

1212
for agent in agents:
1313
agent["worker_group"] = "xx"
14+
agent["status"] = "started"
1415

15-
connections.mongodb_jobs.mrq_agents.insert_many(agents + [{"worker_group": "zz"}])
16+
connections.mongodb_jobs.mrq_agents.insert_many(agents + [{"worker_group": "yy", "status": "started"}, {"worker_group": "zz"}])
1617
connections.mongodb_jobs.mrq_workergroups.insert_one({"_id": "xx", "profiles": profiles})
1718

1819
agent = Agent(worker_group="xx")

tests/test_general.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -152,6 +152,8 @@ def test_general_simple_task_reverse(worker):
152152
def test_known_queues_lifecycle(worker):
153153

154154
worker.start(queues="default_reverse xtest test_timed_set", flags="--config tests/fixtures/config-raw1.py")
155+
time.sleep(1)
156+
worker.wait_for_idle()
155157

156158
# Test known queues
157159
from mrq.queue import Queue, send_task

tests/test_parallel.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -60,7 +60,7 @@ def test_dequeue_strategy(worker, p_greenlets, p_strategy):
6060
print("Worker has flags %s" % flags)
6161
worker.start(flags=flags, queues="q1 q2", deps=False, block=False)
6262

63-
gotit = worker.wait_for_idle()
63+
gotit = worker.wait_for_idle(timeout=10)
6464

6565
if p_strategy == "burst":
6666
assert not gotit # because worker should be stopped already

tests/test_subpool.py

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -49,13 +49,22 @@ def run(params):
4949
assert total_time < 2
5050

5151

52-
def test_subpool_exception(worker):
52+
@pytest.mark.parametrize(["p_imap"], [
53+
[True],
54+
[False]
55+
])
56+
def test_subpool_exception(worker, p_imap):
5357

54-
# Exception
58+
# An exception in the subpool is raised outside the pool
5559
worker.send_task("tests.tasks.general.SubPool", {
56-
"pool_size": 20, "inner_params": ["exception"]
60+
"pool_size": 20, "inner_params": ["exception"], "imap": p_imap
5761
}, accept_statuses=["failed"])
5862

63+
job = worker.mongodb_jobs.mrq_jobs.find_one()
64+
assert job
65+
assert job["status"] == "failed"
66+
assert "__INNER_EXCEPTION_LINE__" in job["traceback"]
67+
5968

6069
@pytest.mark.parametrize(["p_size"], [
6170
[0],

0 commit comments

Comments
 (0)