Merge branch 'release_19.01' into release_19.05

This commit is contained in:
Nicola Soranzo
2019-05-07 16:51:58 +01:00
3 changed files with 19 additions and 49 deletions
+2 -6
View File
@@ -28,13 +28,9 @@ SLURM_MEMORY_LIMIT_EXCEEDED_MSG = 'slurmstepd: error: Exceeded job memory limit'
SLURM_MEMORY_LIMIT_EXCEEDED_PARTIAL_WARNINGS = [': Exceeded job memory limit at some point.',
': Exceeded step memory limit at some point.']
SLURM_MEMORY_LIMIT_SCAN_SIZE = 16 * 1024 * 1024 # 16MB
SLURM_UNABLE_TO_ADD_TASK_TO_MEMORY_CG_MSG_RE = re.compile(r"""slurmstepd: error: task/cgroup: unable to add task\[pid=\d+\] to memory cg '\(null\)'$""")
SLURM_UNABLE_TO_CREATE_CGROUP_MSG_RE = re.compile(r"""slurmstepd: error: xcgroup_instantiate: unable to create cgroup '[^']+' : No space left on device$""")
SLURM_UNABLE_TO_INSTANCIATE_CGROUP_MSG_RE = re.compile(r"""slurmstepd: error: jobacct_gather/cgroup: unable to instanciate (job|user) \d+ memory cgroup$""")
SLURM_CGROUP_RE = re.compile(r"""slurmstepd: .*cgroup.*$""")
SLURM_TOP_WARNING_RES = (
SLURM_UNABLE_TO_ADD_TASK_TO_MEMORY_CG_MSG_RE,
SLURM_UNABLE_TO_CREATE_CGROUP_MSG_RE,
SLURM_UNABLE_TO_INSTANCIATE_CGROUP_MSG_RE
SLURM_CGROUP_RE,
)
# These messages are returned to the user
+8 -39
View File
@@ -5,13 +5,10 @@ This parallelizes the task over available cores using multiprocessing.
Code mostly taken form CloudBioLinux.
"""
import contextlib
import functools
import glob
import multiprocessing
import os
import subprocess
from multiprocessing.pool import IMapIterator
import threading
try:
import boto
@@ -20,13 +17,6 @@ except ImportError:
boto = None
def map_wrap(f):
@functools.wraps(f)
def wrapper(args):
return f(*args)
return wrapper
def mp_from_ids(s3server, mp_id, mp_keyname, mp_bucketname):
"""Get the multipart upload from the bucket and multipart IDs.
@@ -51,7 +41,6 @@ def mp_from_ids(s3server, mp_id, mp_keyname, mp_bucketname):
return mp
@map_wrap
def transfer_part(s3server, mp_id, mp_keyname, mp_bucketname, i, part):
"""Transfer a part of a multipart upload. Designed to be run in parallel.
"""
@@ -64,8 +53,6 @@ def transfer_part(s3server, mp_id, mp_keyname, mp_bucketname, i, part):
def multipart_upload(s3server, bucket, s3_key_name, tarball, mb_size):
"""Upload large files using Amazon's multipart upload functionality.
"""
cores = multiprocessing.cpu_count()
def split_file(in_file, mb_size, split_num=5):
prefix = os.path.join(os.path.dirname(in_file),
"%sS3PART" % (os.path.basename(s3_key_name)))
@@ -80,29 +67,11 @@ def multipart_upload(s3server, bucket, s3_key_name, tarball, mb_size):
mp = bucket.initiate_multipart_upload(s3_key_name,
reduced_redundancy=s3server['use_rr'])
with multimap(cores) as pmap:
for _ in pmap(transfer_part, ((s3server, mp.id, mp.key_name, mp.bucket_name, i, part)
for (i, part) in
enumerate(split_file(tarball, mb_size, cores)))):
pass
for (i, part) in enumerate(split_file(tarball, mb_size)):
t = threading.Thread(
target=transfer_part,
args=(s3server, mp.id, mp.key_name, mp.bucket_name, i, part))
t.start()
t.join()
mp.complete_upload()
@contextlib.contextmanager
def multimap(cores=None):
"""Provide multiprocessing imap like function.
The context manager handles setting up the pool, worked around interrupt issues
and terminating the pool on completion.
"""
if cores is None:
cores = max(multiprocessing.cpu_count() - 1, 1)
def wrapper(func):
def wrap(self, timeout=None):
return func(self, timeout=timeout if timeout is not None else 1e100)
return wrap
IMapIterator.next = wrapper(IMapIterator.next)
pool = multiprocessing.Pool(cores)
yield pool.imap
pool.terminate()
+9 -4
View File
@@ -62,7 +62,8 @@ class ConfiguresHandlers(object):
handler_id,
[x.strip() for x in handler.get('tags', self.DEFAULT_HANDLER_TAG).split(',')]
)
self.default_handler_id = self._get_default(self.app.config, config_element, list(self.handlers.keys()))
self.default_handler_id = self._get_default(
self.app.config, config_element, list(self.handlers.keys()), required=False)
def _init_handler_assignment_methods(self, config_element=None):
self.__is_handler = None
@@ -116,7 +117,7 @@ class ConfiguresHandlers(object):
def _parse_handler(self, handler_id, handler_def):
pass
def _get_default(self, config, parent, names, auto=False):
def _get_default(self, config, parent, names, auto=False, required=True):
"""
Returns the default attribute set in a parent tag like <handlers> or
<destinations>, or return the ID of the child, if there is no explicit
@@ -128,6 +129,8 @@ class ConfiguresHandlers(object):
:type names: list of str
:param auto: Automatically set a default if there is no default in the parent tag and there is only one child.
:type auto: bool
:param required: Require a default to be set or determined automatically, else raise Exception
:type required: bool
:returns: str -- id or tag representing the default.
"""
@@ -142,12 +145,14 @@ class ConfiguresHandlers(object):
if rval is not None:
# If the parent element has a 'default' attribute, use the id or tag in that attribute
if self.deterministic_handler_assignment and rval not in names:
if required and rval not in names:
raise Exception("<%s> default attribute '%s' does not match a defined id or tag in a child element" % (parent.tag, rval))
log.debug("<%s> default set to child with id or tag '%s'" % (parent.tag, rval))
elif auto and len(names) == 1:
log.info("Setting <%s> default to child with id '%s'" % (parent.tag, names[0]))
rval = names[0]
elif required:
raise Exception("No <%s> default specified, please specify a valid id or tag with the 'default' attribute" % parent.tag)
return rval
def _findall_with_required(self, parent, match, attribs=None):
@@ -181,7 +186,7 @@ class ConfiguresHandlers(object):
@property
def deterministic_handler_assignment(self):
return self.handler_assignment_methods and all(
return self.handler_assignment_methods and any(
filter(lambda x: x in (
HANDLER_ASSIGNMENT_METHODS.UWSGI_MULE_MESSAGE,
HANDLER_ASSIGNMENT_METHODS.DB_PREASSIGN,