mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
Merge pull request #6924 from natefoo/quiesce-handler
Be more precise with the job runner worker thread join timeout and improve documentation
This commit is contained in:
@@ -1019,9 +1019,14 @@ galaxy:
|
||||
|
||||
# When stopping Galaxy cleanly, how much time to give various
|
||||
# monitoring/polling threads to finish before giving up on joining
|
||||
# them. Set to 0 to disable this and restore the pre-18.01 default
|
||||
# behavior.
|
||||
#monitor_thread_join_timeout: 5
|
||||
# them. Set to 0 to disable this and terminate without waiting. Among
|
||||
# others, these threads include the job handler workers, which are
|
||||
# responsible for preparing/submitting and collecting/finishing jobs,
|
||||
# and which can cause job errors if not shut down cleanly. If using
|
||||
# supervisord, consider also increasing the value of `stopwaitsecs`.
|
||||
# If using job handler mules, consider also setting the `mule-reload-
|
||||
# mercy` uWSGI option. See the Galaxy Admin Documentation for more.
|
||||
#monitor_thread_join_timeout: 30
|
||||
|
||||
# Write thread status periodically to 'heartbeat.log', (careful, uses
|
||||
# disk space rapidly!). Useful to determine why your processes may be
|
||||
|
||||
@@ -2097,9 +2097,15 @@
|
||||
:Description:
|
||||
When stopping Galaxy cleanly, how much time to give various
|
||||
monitoring/polling threads to finish before giving up on joining
|
||||
them. Set to 0 to disable this and restore the pre-18.01 default
|
||||
behavior.
|
||||
:Default: ``5``
|
||||
them. Set to 0 to disable this and terminate without waiting.
|
||||
Among others, these threads include the job handler workers, which
|
||||
are responsible for preparing/submitting and collecting/finishing
|
||||
jobs, and which can cause job errors if not shut down cleanly. If
|
||||
using supervisord, consider also increasing the value of
|
||||
`stopwaitsecs`. If using job handler mules, consider also setting
|
||||
the `mule-reload-mercy` uWSGI option. See the Galaxy Admin
|
||||
Documentation for more.
|
||||
:Default: ``30``
|
||||
:Type: int
|
||||
|
||||
|
||||
|
||||
@@ -407,6 +407,22 @@ option, `enable-threads` is set implicitly. This option enables the Python GIL a
|
||||
by Galaxy itself for various non-web tasks), which Galaxy uses extensively. Setting it explicitly, however, is harmless
|
||||
and can prevent strange difficult-to-debug situations if `threads` is accidentally unset.
|
||||
|
||||
**Worker/Mule shutdown/reload mercy**
|
||||
|
||||
By default, uWSGI will wait up to 60 seconds for web workers and mules to terminate. This is generally safe for
|
||||
servicing web requests, but some parts of Galaxy's job preparation/submission and collection/finishing operations can
|
||||
take quite a bit of time to complete and are not entirely reentrant: job errors or state inconsistencies can occur if
|
||||
interrupted (although every effort has been made to minimize such possibilities). By default, Galaxy will wait up to 30
|
||||
seconds for the threads allocated for these operations to terminate after instructing them to shut down. You can change
|
||||
this behavior by increasing the value of `monitor_thread_join_timeout` in the `galaxy` section of `galaxy.yml`. The
|
||||
maximum amount of time that Galaxy will take to shut down job runner workers is `monitor_thread_join_timeout *
|
||||
runner_plugin_count` since each plugin is shut down sequentially (`runner_plugin_count` is the number of `<plugin>`s in
|
||||
your `job_conf.xml`).
|
||||
|
||||
Thus you should set the appropriate uWSGI `*-restart-mercy` option to a value higher than the maximum job runner worker
|
||||
shutdown time. If using **uWSGI all-in-one**, set `worker-reload-mercy`, and if using **uWSGI + Mule job handling**, set
|
||||
`mule-reload-mercy` (both in the `uwsgi` section of `galaxy.yml`).
|
||||
|
||||
**Signals**
|
||||
|
||||
The signal handling options (`die-on-term` and `hook-master-start` with `unix_signal` values) are not required but, if
|
||||
@@ -494,6 +510,7 @@ umask = 022
|
||||
autostart = true
|
||||
autorestart = true
|
||||
startsecs = 10
|
||||
stopwaitsecs = 65
|
||||
user = galaxy
|
||||
numprocs = 1
|
||||
stopsignal = INT
|
||||
@@ -506,6 +523,9 @@ must stay up for this long before we consider it OK. If the process crashes soon
|
||||
made to your local installation) supervisord will try again a couple of times to restart the process before giving up
|
||||
and marking it as failed. This is one of the many ways supervisord is much friendly for managing these sorts of tasks.
|
||||
|
||||
The value of `stopwaitsecs` should be at least as large as the smallest value of uWSGI's `reload-mercy`,
|
||||
`worker-reload-mercy`, and `mule-reload-mercy` options, all of which default to `60`.
|
||||
|
||||
If using the **uWSGI + Webless** scenario, you'll need to addtionally define job handlers to start. There's no simple
|
||||
way to activate a virtualenv when using supervisor, but you can simulate the effects by setting `$PATH` and
|
||||
`$VIRTUAL_ENV`:
|
||||
@@ -520,6 +540,7 @@ umask = 022
|
||||
autostart = true
|
||||
autorestart = true
|
||||
startsecs = 15
|
||||
stopwaitsecs = 35
|
||||
user = galaxy
|
||||
environment = VIRTUAL_ENV="/srv/galaxy/venv",PATH="/srv/galaxy/venv/bin:%(ENV_PATH)s"
|
||||
```
|
||||
@@ -529,6 +550,10 @@ substitution in the `command` and `process_name` fields. We've set `numprocs = 3
|
||||
processes. Supervisord will loop over `0..numprocs` and launch `handler0`, `handler1`, and `handler2` processes
|
||||
automatically for us, templating out the command string so each handler receives a different log file and name.
|
||||
|
||||
The value of `stopwaitsecs` should be at least as large as `monitor_thread_join_timeout * runner_plugin_count`, which
|
||||
is `30` in the default configuration (`monitor_thread_join_timeout` is a Galaxy configuration option and
|
||||
`runner_plugin_count` is the number of `<plugin>`s in your `job_conf.xml`).
|
||||
|
||||
Lastly, collect the tasks defined above into a single group. If you are not using webless handlers this is as simple as:
|
||||
|
||||
```ini
|
||||
|
||||
@@ -425,7 +425,7 @@ class Configuration(object):
|
||||
self.use_heartbeat = string_as_bool(kwargs.get('use_heartbeat', 'False'))
|
||||
self.heartbeat_interval = int(kwargs.get('heartbeat_interval', 20))
|
||||
self.heartbeat_log = kwargs.get('heartbeat_log', None)
|
||||
self.monitor_thread_join_timeout = int(kwargs.get("monitor_thread_join_timeout", 5))
|
||||
self.monitor_thread_join_timeout = int(kwargs.get("monitor_thread_join_timeout", 30))
|
||||
self.log_actions = string_as_bool(kwargs.get('log_actions', 'False'))
|
||||
self.log_events = string_as_bool(kwargs.get('log_events', 'False'))
|
||||
self.sanitize_all_html = string_as_bool(kwargs.get('sanitize_all_html', True))
|
||||
|
||||
@@ -965,8 +965,12 @@ class DefaultJobDispatcher(object):
|
||||
job_wrapper.fail(DEFAULT_JOB_PUT_FAILURE_MESSAGE)
|
||||
|
||||
def shutdown(self):
|
||||
for runner in self.job_runners.values():
|
||||
failures = []
|
||||
for name, runner in self.job_runners.items():
|
||||
try:
|
||||
runner.shutdown()
|
||||
except Exception:
|
||||
raise Exception("Failed to shutdown runner %s" % runner)
|
||||
failures.append(name)
|
||||
log.exception("Failed to shutdown runner %s", name)
|
||||
if failures:
|
||||
raise Exception("Failed to shutdown runners: %s" % ', '.join(failures))
|
||||
|
||||
@@ -6,8 +6,10 @@ import logging
|
||||
import os
|
||||
import string
|
||||
import subprocess
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
import traceback
|
||||
|
||||
from six.moves.queue import (
|
||||
Empty,
|
||||
@@ -85,10 +87,21 @@ class BaseJobRunner(object):
|
||||
log.debug('Starting %s %s workers' % (self.nworkers, self.runner_name))
|
||||
for i in range(self.nworkers):
|
||||
worker = threading.Thread(name="%s.work_thread-%d" % (self.runner_name, i), target=self.run_next)
|
||||
worker.setDaemon(True)
|
||||
worker.daemon = True
|
||||
self.app.application_stack.register_postfork_function(worker.start)
|
||||
self.work_threads.append(worker)
|
||||
|
||||
def _alive_worker_threads(self, cycle=False):
|
||||
# yield endlessly as long as there are alive threads if cycle is True
|
||||
alive = True
|
||||
while alive:
|
||||
alive = False
|
||||
for thread in self.work_threads:
|
||||
if thread.is_alive():
|
||||
if cycle:
|
||||
alive = True
|
||||
yield thread
|
||||
|
||||
def run_next(self):
|
||||
"""Run the next item in the work queue (a job waiting to run)
|
||||
"""
|
||||
@@ -129,22 +142,36 @@ class BaseJobRunner(object):
|
||||
def shutdown(self):
|
||||
"""Attempts to gracefully shut down the worker threads
|
||||
"""
|
||||
log.info("%s: Sending stop signal to %s worker threads" % (self.runner_name, len(self.work_threads)))
|
||||
log.info("%s: Sending stop signal to %s job worker threads", self.runner_name, len(self.work_threads))
|
||||
for i in range(len(self.work_threads)):
|
||||
self.work_queue.put((STOP_SIGNAL, None))
|
||||
|
||||
join_timeout = self.app.config.monitor_thread_join_timeout
|
||||
if join_timeout > 0:
|
||||
exception = None
|
||||
for thread in self.work_threads:
|
||||
log.info("Waiting up to %d seconds for job worker threads to shutdown...", join_timeout)
|
||||
start = time.time()
|
||||
# NOTE: threads that have already joined by now are not going to be logged
|
||||
for thread in self._alive_worker_threads(cycle=True):
|
||||
if time.time() > (start + join_timeout):
|
||||
break
|
||||
try:
|
||||
thread.join(join_timeout)
|
||||
except Exception as e:
|
||||
exception = e
|
||||
log.exception("Faild to shutdown worker thread")
|
||||
thread.join(2)
|
||||
except Exception:
|
||||
log.exception("Caught exception attempting to shutdown job worker thread %s:", thread.name)
|
||||
if not thread.is_alive():
|
||||
log.debug("Job worker thread terminated: %s", thread.name)
|
||||
else:
|
||||
log.info("All job worker threads shutdown cleanly")
|
||||
return
|
||||
|
||||
if exception:
|
||||
raise exception
|
||||
for thread in self._alive_worker_threads():
|
||||
try:
|
||||
frame = sys._current_frames()[thread.ident]
|
||||
except KeyError:
|
||||
# thread is now stopped
|
||||
continue
|
||||
log.warning("Timed out waiting for job worker thread %s to terminate, shutdown will be unclean! Thread "
|
||||
"stack is:\n%s", thread.name, ''.join(traceback.format_stack(frame)))
|
||||
|
||||
# Most runners should override the legacy URL handler methods and destination param method
|
||||
def url_to_destination(self, url):
|
||||
|
||||
@@ -1563,12 +1563,17 @@ mapping:
|
||||
|
||||
monitor_thread_join_timeout:
|
||||
type: int
|
||||
default: 5
|
||||
default: 30
|
||||
required: false
|
||||
desc: |
|
||||
When stopping Galaxy cleanly, how much time to give various monitoring/polling
|
||||
threads to finish before giving up on joining them. Set to 0 to disable this and
|
||||
restore the pre-18.01 default behavior.
|
||||
terminate without waiting. Among others, these threads include the job handler
|
||||
workers, which are responsible for preparing/submitting and collecting/finishing
|
||||
jobs, and which can cause job errors if not shut down cleanly. If using
|
||||
supervisord, consider also increasing the value of `stopwaitsecs`. If using job
|
||||
handler mules, consider also setting the `mule-reload-mercy` uWSGI option. See
|
||||
the Galaxy Admin Documentation for more.
|
||||
|
||||
use_heartbeat:
|
||||
type: bool
|
||||
|
||||
@@ -228,6 +228,7 @@ def setup_galaxy_config(
|
||||
user_library_import_dir=user_library_import_dir,
|
||||
webhooks_dir=TEST_WEBHOOKS_DIR,
|
||||
logging=LOGGING_CONFIG_DEFAULT,
|
||||
monitor_thread_join_timeout=5,
|
||||
)
|
||||
config.update(database_conf(tmpdir, prefer_template_database=prefer_template_database))
|
||||
config.update(install_database_conf(tmpdir, default_merged=default_install_db_merged))
|
||||
|
||||
Reference in New Issue
Block a user