Skip to content

Commit 7cc558c

Browse files
committed
Merge branch 'master' of github.com:pricingassistant/mrq
* 'master' of github.com:pricingassistant/mrq: Added a --max_memory option to restart workers when they go OOM
2 parents e46726a + 721d891 commit 7cc558c

3 files changed

Lines changed: 48 additions & 0 deletions

File tree

mrq/config.py

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -235,6 +235,14 @@ def add_parser_args(parser, config_type):
235235
help='Gevent: max number of jobs to do before quitting.' +
236236
' Temp workaround for memory leaks')
237237

238+
parser.add_argument(
239+
'--max_memory',
240+
default=0,
241+
type=int,
242+
action='store',
243+
help='Max memory (in Mb) after which the process will be shut down. Use with --processes [1-N]' +
244+
'to have supervisord automatically respawn the worker when this happens')
245+
238246
parser.add_argument(
239247
'--greenlets',
240248
'--gevent', # deprecated

mrq/worker.py

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -300,6 +300,10 @@ def report_worker(self, w=0):
300300

301301
report = self.get_worker_report(with_memory=True)
302302

303+
if self.config["max_memory"] > 0:
304+
if report["process"]["mem"]["total"] > (self.config["max_memory"] * 1024 * 1024):
305+
self.shutdown_max_memory()
306+
303307
if self.config["report_file"]:
304308
with open(self.config["report_file"], "wb") as f:
305309
f.write(json.dumps(report, ensure_ascii=False)) # pylint: disable=no-member
@@ -398,6 +402,9 @@ def work_loop(self, max_jobs=None):
398402

399403
while True:
400404

405+
if self.graceful_stop:
406+
break
407+
401408
while True:
402409

403410
free_pool_slots = self.gevent_pool.free_count()
@@ -567,6 +574,14 @@ def shutdown_graceful(self):
567574
self.log.info("Graceful shutdown...")
568575
raise StopRequested() # pylint: disable=nonstandard-exception
569576

577+
def shutdown_max_memory(self):
578+
579+
# Not in the exitcodes list: we want it to be restarted.
580+
self.exitcode = 4
581+
582+
self.log.info("Max memory reached, shutdown...")
583+
self.graceful_stop = True
584+
570585
def shutdown_now(self):
571586
""" Forced shutdown: interrupts all the jobs. """
572587

tests/test_memoryleaks.py

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,27 @@
11
import time
22

33

4+
def test_max_memory_restart(worker):
5+
6+
N = 20
7+
8+
worker.start(
9+
flags="--processes 1 --greenlets 1 --max_memory 50 --report_interval 1")
10+
11+
worker.send_tasks(
12+
"tests.tasks.general.Leak",
13+
[{"size": 1000000, "sleep": 1} for _ in range(N)],
14+
queue="default",
15+
block=True
16+
)
17+
18+
assert worker.mongodb_jobs.mrq_jobs.find(
19+
{"status": "success"}).count() == N
20+
21+
# We must have been restarted at least once.
22+
assert worker.mongodb_jobs.mrq_workers.find().count() > 1
23+
24+
425
def get_diff_after_jobs(worker, n_tasks, leak, sleep=0):
526

627
time.sleep(3)
@@ -51,6 +72,8 @@ def test_memoryleaks_noleak(worker):
5172
assert diff100 < 10000
5273
assert diff200 < 10000
5374

75+
assert worker.mongodb_jobs.mrq_workers.find().count() == 1
76+
5477

5578
def test_memoryleaks_1mleak(worker):
5679

@@ -70,3 +93,5 @@ def test_memoryleaks_1mleak(worker):
7093

7194
assert worker.mongodb_jobs.mrq_jobs.find(
7295
{"memory_diff": {"$gte": 80000}}).count() == 10
96+
97+
assert worker.mongodb_jobs.mrq_workers.find().count() == 1

0 commit comments

Comments
 (0)