【问题标题】:_multiprocessing.SemLock is not implemented when running on AWS Lambda在 AWS Lambda 上运行时未实现 _multiprocessing.SemLock
【发布时间】:2016-03-04 12:30:51
【问题描述】:

我有一个使用multiprocessing 包的短代码,并且在我的本地机器上运行良好。

当我上传到AWS Lambda 并在那里运行时,我收到以下错误(stacktrace 已修整):

[Errno 38] Function not implemented: OSError
Traceback (most recent call last):
  File "/var/task/recorder.py", line 41, in record
    pool = multiprocessing.Pool(10)
  File "/usr/lib64/python2.7/multiprocessing/__init__.py", line 232, in Pool
    return Pool(processes, initializer, initargs, maxtasksperchild)
  File "/usr/lib64/python2.7/multiprocessing/pool.py", line 138, in __init__
    self._setup_queues()
  File "/usr/lib64/python2.7/multiprocessing/pool.py", line 234, in _setup_queues
    self._inqueue = SimpleQueue()
  File "/usr/lib64/python2.7/multiprocessing/queues.py", line 354, in __init__
    self._rlock = Lock()
  File "/usr/lib64/python2.7/multiprocessing/synchronize.py", line 147, in __init__
    SemLock.__init__(self, SEMAPHORE, 1, 1)
  File "/usr/lib64/python2.7/multiprocessing/synchronize.py", line 75, in __init__
    sl = self._semlock = _multiprocessing.SemLock(kind, value, maxvalue)
OSError: [Errno 38] Function not implemented

会不会是python核心包的一部分没有实现?我不知道我在下面运行什么,所以我无法在那里登录和调试。

任何想法如何在 Lambda 上运行 multiprocessing?

【问题讨论】:

标签: python-multiprocessing aws-lambda


【解决方案1】:

据我所知,多处理无法在 AWS Lambda 上运行,因为缺少执行环境/容器/dev/shm - 请参阅https://forums.aws.amazon.com/thread.jspa?threadID=219962(可能需要登录)。

没有关于亚马逊是否/何时会改变这一点的消息(我能找到)。我还查看了其他库,例如如果找不到/dev/shm、but that doesn't actually solve the problem,https://pythonhosted.org/joblib/parallel.html 将回退到/tmp(我们知道确实存在)。

【讨论】:

  • 你能详细说明如何用joblib解决这个问题吗?我现在正在对其进行测试,并且 joblib 无法恢复到串行操作:[Errno 38] Function not implemented. joblib will operate in serial mode
  • This thread 似乎表明 joblib 实际上无法解决此问题。
  • 是的,抱歉,我从未对此进行过深入研究。很可能是行不通的。
  • 请更新您的答案。在访问者必须阅读 cmets 之前,它看起来具有误导性。
【解决方案2】:

multiprocessing.Pool 和 multiprocessing.Queue 本身不受支持(因为 SemLock 的问题),但 multiprocessing.Process 和 multiprocessing.Pipe 等在 AWSLambda 中可以正常工作。

这应该允许您通过手动创建/派生进程并使用multiprocessing.Pipe 在父进程和子进程之间进行通信来构建解决方案。希望有帮助

【讨论】:

  • multiprocessing.Queue 对我不起作用,我得到与问题中相同的错误。
  • 队列不起作用,没有 /dev/shm 就不能在进程之间做任何锁
【解决方案3】:

您可以使用 Python 的多处理模块在 AWS Lambda 上并行运行例程,但您不能使用其他答案中所述的池或队列。一个可行的解决方案是使用本文中概述的 Process 和 Pipe https://aws.amazon.com/blogs/compute/parallel-processing-in-python-with-aws-lambda/

虽然这篇文章确实帮助我找到了解决方案(在下面分享),但仍有一些事情需要注意。首先,基于 Process 和 Pipe 的解决方案不如 Pool 中内置的 map 函数快,尽管我确实看到了几乎线性的加速,因为我增加了 Lambda 函数中的可用内存/CPU 资源。其次,以这种方式开发多处理功能时,必须进行相当多的管理。我怀疑这至少是我的解决方案比内置方法慢的部分原因。如果有人有加快速度的建议,我很乐意听到他们的声音!最后,虽然文章指出多处理对于卸载异步进程很有用,但使用多处理还有其他原因,例如我正在尝试做的大量密集数学运算。最后,我对性能提升感到非常满意,因为它比顺序执行要好得多!

代码:

# Python 3.6
from multiprocessing import Pipe, Process

def myWorkFunc(data, connection):
    result = None

    # Do some work and store it in result

    if result:
        connection.send([result])
    else:
        connection.send([None])


def myPipedMultiProcessFunc():

    # Get number of available logical cores
    plimit = multiprocessing.cpu_count()

    # Setup management variables
    results = []
    parent_conns = []
    processes = []
    pcount = 0
    pactive = []
    i = 0

    for data in iterable:
        # Create the pipe for parent-child process communication
        parent_conn, child_conn = Pipe()
        # create the process, pass data to be operated on and connection
        process = Process(target=myWorkFunc, args=(data, child_conn,))
        parent_conns.append(parent_conn)
        process.start()
        pcount += 1

        if pcount == plimit: # There is not currently room for another process
            # Wait until there are results in the Pipes
            finishedConns = multiprocessing.connection.wait(parent_conns)
            # Collect the results and remove the connection as processing
            # the connection again will lead to errors
            for conn in finishedConns:
                results.append(conn.recv()[0])
                parent_conns.remove(conn)
                # Decrement pcount so we can add a new process
                pcount -= 1

    # Ensure all remaining active processes have their results collected
    for conn in parent_conns:
        results.append(conn.recv()[0])
        conn.close()

    # Process results as needed

【讨论】:

  • 这段代码有点难以理解。 myPipedMultiProcessFunc 是 Pool.map() 的可行替代品吗?
  • 如myPipedMultiProcessFunc 所写,比在顺序循环中运行myWorkFunc 快得多。写这篇文章已经有一段时间了,但我记得这个实现大约是Pool.map() 速度的 80%。如果我的代码中有一些不清楚的地方,很高兴跟进。
【解决方案4】:

我遇到了同样的问题。这是我之前在本地机器上运行良好的代码:

import concurrent.futures


class Concurrent:

    @staticmethod
    def execute_concurrently(function, kwargs_list):
        results = []
        with concurrent.futures.ProcessPoolExecutor() as executor:
            for _, result in zip(kwargs_list, executor.map(function, kwargs_list)):
                results.append(result)
        return results

我用这个替换了它:

import concurrent.futures


class Concurrent:

    @staticmethod
    def execute_concurrently(function, kwargs_list):
        results = []
        with concurrent.futures.ThreadPoolExecutor(max_workers=10) as executor:
            futures = [executor.submit(function, kwargs) for kwargs in kwargs_list]
        for future in concurrent.futures.as_completed(futures):
            results.append(future.result())
        return results

像魅力一样工作。

取自this pull request

【讨论】:

  • 请注意,这会将其更改为使用多线程而不是多处理。它会运行,但取决于每个函数执行的功能,其性能可能不如使用多处理。
【解决方案5】:

如果您可以在 AWS Lambda 上使用 Python 3.7(或更早版本),则应该没问题,因为它不使用 SemLock。

但是,如果您只需要 AWS Lambda 上的 async_results(没有任何额外要求)和更高版本的 Python,这里有一个更新的直接替换(基于 https://code.activestate.com/recipes/576519-thread-pool-with-same-api-as-multiprocessingpool/):

import sys
import threading
from queue import Empty, Queue

SENTINEL = "QUIT"

def is_sentinel(obj):
    """
    Predicate to determine whether an item from the queue is the
    signal to stop
    """
    return type(obj) is str and obj == SENTINEL

class TimeoutError(Exception):
    """
    Raised when a result is not available within the given timeout
    """

class Pool(object):
    def __init__(self, processes, name="Pool"):
        self.processes = processes
        self._queue = Queue()
        self._closed = False
        self._workers = []
        for idx in range(processes):
            thread = PoolWorker(self._queue, name="Worker-%s-%d" % (name, idx))
            try:
                thread.start()
            except Exception:
                # If one thread has a problem, undo everything
                self.terminate()
                raise
            else:
                self._workers.append(thread)

    def apply_async(self, func, args, kwds):
        apply_result = ApplyResult()
        job = Job(func, args, kwds, apply_result)
        self._queue.put(job)
        return apply_result

    def close(self):
        self._closed = True

    def join(self):
        """
        This is only called when all are done.
        """
        self.terminate()

    def terminate(self):
        """
        Stops the worker processes immediately without completing
        outstanding work. When the pool object is garbage collected
        terminate() will be called immediately.
        """
        self.close()

        # Clearing the job queue
        try:
            while True:
                self._queue.get_nowait()
        except Empty:
            pass

        for thread in self._workers:
            self._queue.put(SENTINEL)


class PoolWorker(threading.Thread):
    """
    Thread that consumes WorkUnits from a queue to process them
    """
    def __init__(self, queue, *args, **kwds):
        """
        Args:
            queue: the queue of jobs
        """
        threading.Thread.__init__(self, *args, **kwds)
        self.daemon = True
        self._queue = queue

    def run(self):
        """
        Process the job, or wait for sentinel to exit
        """
        while True:
            job = self._queue.get()
            if is_sentinel(job):
                # Got sentinel
                break
            job.process()


class ApplyResult(object):
    """
    Container to hold results.
    """
    def __init__(self):
        self._data = None
        self._success = None
        self._event = threading.Event()

    def ready(self):
        is_ready = self._event.isSet()
        return is_ready

    def get(self, timeout=None):
        """
        Returns the result when it arrives. If timeout is not None and
        the result does not arrive within timeout seconds then
        TimeoutError is raised. If the remote call raised an exception
        then that exception will be reraised by get().
        """
        if not self.wait(timeout):
            raise TimeoutError("Result not available within %fs" % timeout)
        if self._success:
            return self._data
        raise self._data[0](self._data[1], self._data[2])

    def wait(self, timeout=None):
        """
        Waits until the result is available or until timeout
        seconds pass.
        """
        self._event.wait(timeout)
        return self._event.isSet()

    def _set_exception(self):
        self._data = sys.exc_info()
        self._success = False
        self._event.set()

    def _set_value(self, value):
        self._data = value
        self._success = True
        self._event.set()


class Job(object):
    """
    A work unit that corresponds to the execution of a single function
    """
    def __init__(self, func, args, kwds, apply_result):
        """
        Args:
            func: function
            args: function args
            kwds: function kwargs
            apply_result: ApplyResult object that holds the result
                of the function call
        """
        self._func = func
        self._args = args
        self._kwds = kwds
        self._result = apply_result

    def process(self):
        """
        Call the function with the args/kwds and tell the ApplyResult
        that its result is ready. Correctly handles the exceptions
        happening during the execution of the function
        """
        try:
            result = self._func(*self._args, **self._kwds)
        except Exception:
            self._result._set_exception()
        else:
            self._result._set_value(result)

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-04-07
    • 1970-01-01
    • 2021-02-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-08-14
    • 1970-01-01
    相关资源
    最近更新 更多