Merge pull request #2269 from bmiklautz/spawn_v2
Fix scheduler issue on linux related to multiprocessing
This commit is contained in:
+16
-6
@@ -596,12 +596,12 @@ class Scheduler(threading.Thread):
|
|||||||
utc_time(bool): do all datetime calculations assuming UTC as the
|
utc_time(bool): do all datetime calculations assuming UTC as the
|
||||||
timezone. Remember to pass `start_time` and `stop_time` to tasks
|
timezone. Remember to pass `start_time` and `stop_time` to tasks
|
||||||
accordingly
|
accordingly
|
||||||
|
use_spawn(bool): use spawn for subprocess (only useable with python3)
|
||||||
"""
|
"""
|
||||||
|
|
||||||
def __init__(self, db, tasks=None, migrate=True,
|
def __init__(self, db, tasks=None, migrate=True,
|
||||||
worker_name=None, group_names=None, heartbeat=HEARTBEAT,
|
worker_name=None, group_names=None, heartbeat=HEARTBEAT,
|
||||||
max_empty_runs=0, discard_results=False, utc_time=False):
|
max_empty_runs=0, discard_results=False, utc_time=False, use_spawn=False):
|
||||||
|
|
||||||
threading.Thread.__init__(self)
|
threading.Thread.__init__(self)
|
||||||
self.setDaemon(True)
|
self.setDaemon(True)
|
||||||
@@ -639,6 +639,7 @@ class Scheduler(threading.Thread):
|
|||||||
current._scheduler = self
|
current._scheduler = self
|
||||||
|
|
||||||
self.define_tables(db, migrate=migrate)
|
self.define_tables(db, migrate=migrate)
|
||||||
|
self.use_spawn = use_spawn
|
||||||
|
|
||||||
def execute(self, task):
|
def execute(self, task):
|
||||||
"""Start the background process.
|
"""Start the background process.
|
||||||
@@ -649,10 +650,19 @@ class Scheduler(threading.Thread):
|
|||||||
Returns:
|
Returns:
|
||||||
a `TaskReport` object
|
a `TaskReport` object
|
||||||
"""
|
"""
|
||||||
outq = multiprocessing.Queue()
|
outq = None
|
||||||
retq = multiprocessing.Queue(maxsize=1)
|
retq = None
|
||||||
self.process = p = \
|
if (self.use_spawn and not PY2):
|
||||||
multiprocessing.Process(target=executor, args=(retq, task, outq))
|
ctx = multiprocessing.get_context('spawn')
|
||||||
|
outq = ctx.Queue()
|
||||||
|
retq = ctx.Queue(maxsize=1)
|
||||||
|
sel.process = p = ctx.Process(target=executor, args=(retq, task, outq))
|
||||||
|
else:
|
||||||
|
outq = multiprocessing.Queue()
|
||||||
|
retq = multiprocessing.Queue(maxsize=1)
|
||||||
|
self.process = p = \
|
||||||
|
multiprocessing.Process(target=executor, args=(retq, task, outq))
|
||||||
|
|
||||||
self.process_queues = (retq, outq)
|
self.process_queues = (retq, outq)
|
||||||
|
|
||||||
logger.debug(' task starting')
|
logger.debug(' task starting')
|
||||||
|
|||||||
Reference in New Issue
Block a user