【问题标题】:Asynchronous method call in Python?Python中的异步方法调用?
【发布时间】:2010-11-17 08:44:46
【问题描述】:

我想知道Python 中是否有任何用于异步方法调用的库。如果你能做类似的事情会很棒

@async
def longComputation():
    <code>


token = longComputation()
token.registerCallback(callback_function)
# alternative, polling
while not token.finished():
    doSomethingElse()
    if token.finished():
        result = token.result()

或者异步调用非异步例程

def longComputation()
    <code>

token = asynccall(longComputation())

如果在语言核心中有一个更精细的策略作为本地语言,那就太好了。考虑过吗?

【问题讨论】:

标签: python asynchronous


【解决方案1】:

它不在语言核心中,而是一个非常成熟的库,可以满足您的需求,Twisted。它引入了 Deferred 对象,您可以将回调或错误处理程序(“errbacks”)附加到该对象。 Deferred 基本上是一个函数最终会有结果的“承诺”。

【讨论】:

【解决方案2】:

类似:

import threading

thr = threading.Thread(target=foo, args=(), kwargs={})
thr.start() # Will run "foo"
....
thr.is_alive() # Will return whether foo is running currently
....
thr.join() # Will wait till "foo" is done

有关详细信息,请参阅https://docs.python.org/library/threading.html 的文档。

【讨论】:

  • 是的,如果你只需要异步做事,为什么不直接使用线程呢?毕竟线程比进程轻
  • 重要提示:由于“全局解释器锁”,线程的标准实现 (CPython) 对计算绑定任务没有帮助。见图书馆文档:link
  • 使用 thread.join() 真的是异步的吗?如果您不想阻塞线程(例如 UI 线程)并且不想使用大量资源在其上执行 while 循环怎么办?
  • @Mgamerz 加入是同步的。您可以让线程将执行结果放在某个队列中,或/并调用回调。否则你不知道它什么时候完成(如果有的话)。
  • 是否可以像使用 multiprocessing.Pool 一样在线程执行结束时调用回调函数
【解决方案3】:

有什么理由不使用线程吗?您可以使用threading 类。 使用isAlive() 代替finished() 函数。 result() 函数可以join() 线程并检索结果。并且,如果可以的话,重写run() 和__init__ 函数以调用构造函数中指定的函数并将值保存到类的实例中。

【讨论】:

  • 如果它是一个计算量大的函数,线程不会为您带来任何好处(实际上它可能会使事情变慢),因为由于 GIL,Python 进程仅限于一个 CPU 内核。
  • @Kurt,虽然这是真的,但 OP 没有提到性能是他关心的问题。想要异步行为还有其他原因......
  • 当您想选择终止异步方法调用时,python 中的线程并不是很好,因为只有 python 中的主线程接收信号。
【解决方案4】:

您可以使用 Python 2.6 中添加的multiprocessing module。您可以使用进程池,然后通过以下方式异步获取结果:

apply_async(func[, args[, kwds[, callback]]])

例如:

from multiprocessing import Pool

def f(x):
    return x*x

if __name__ == '__main__':
    pool = Pool(processes=1)              # Start a worker processes.
    result = pool.apply_async(f, [10], callback) # Evaluate "f(10)" asynchronously calling callback when finished.

这只是一种选择。这个模块提供了很多功能来实现你想要的。用它来做一个装饰器也很容易。

【讨论】:

  • Lucas S.,不幸的是,您的示例不起作用。回调函数永远不会被调用。
  • 可能值得记住的是,这会产生单独的进程,而不是进程中的单独线程。这可能会产生一些影响。
  • 这有效:result = pool.apply_async(f, [10], callback=finish)
  • 要真正在 python 中异步执行任何操作,需要使用多处理模块来生成新进程。仅仅创建新线程仍然受制于全局解释器锁,它可以防止 python 进程同时执行多项操作。
  • 如果您不想在使用此解决方案时生成新进程 - 将导入更改为 from multiprocessing.dummy import Pool。 multiprocessing.dummy 通过线程而不是进程实现完全相同的行为
【解决方案5】:

我的解决办法是:

import threading

class TimeoutError(RuntimeError):
    pass

class AsyncCall(object):
    def __init__(self, fnc, callback = None):
        self.Callable = fnc
        self.Callback = callback

    def __call__(self, *args, **kwargs):
        self.Thread = threading.Thread(target = self.run, name = self.Callable.__name__, args = args, kwargs = kwargs)
        self.Thread.start()
        return self

    def wait(self, timeout = None):
        self.Thread.join(timeout)
        if self.Thread.isAlive():
            raise TimeoutError()
        else:
            return self.Result

    def run(self, *args, **kwargs):
        self.Result = self.Callable(*args, **kwargs)
        if self.Callback:
            self.Callback(self.Result)

class AsyncMethod(object):
    def __init__(self, fnc, callback=None):
        self.Callable = fnc
        self.Callback = callback

    def __call__(self, *args, **kwargs):
        return AsyncCall(self.Callable, self.Callback)(*args, **kwargs)

def Async(fnc = None, callback = None):
    if fnc == None:
        def AddAsyncCallback(fnc):
            return AsyncMethod(fnc, callback)
        return AddAsyncCallback
    else:
        return AsyncMethod(fnc, callback)

并且完全按照要求工作:

@Async
def fnc():
    pass

【讨论】:

    【解决方案6】:

    您可以实现一个装饰器来使您的函数异步,尽管这有点棘手。 multiprocessing 模块充满了小怪癖和看似随意的限制——不过,更有理由将它封装在一个友好的界面后面。

    from inspect import getmodule
    from multiprocessing import Pool
    
    
    def async(decorated):
        r'''Wraps a top-level function around an asynchronous dispatcher.
    
            when the decorated function is called, a task is submitted to a
            process pool, and a future object is returned, providing access to an
            eventual return value.
    
            The future object has a blocking get() method to access the task
            result: it will return immediately if the job is already done, or block
            until it completes.
    
            This decorator won't work on methods, due to limitations in Python's
            pickling machinery (in principle methods could be made pickleable, but
            good luck on that).
        '''
        # Keeps the original function visible from the module global namespace,
        # under a name consistent to its __name__ attribute. This is necessary for
        # the multiprocessing pickling machinery to work properly.
        module = getmodule(decorated)
        decorated.__name__ += '_original'
        setattr(module, decorated.__name__, decorated)
    
        def send(*args, **opts):
            return async.pool.apply_async(decorated, args, opts)
    
        return send
    

    下面的代码说明了装饰器的用法:

    @async
    def printsum(uid, values):
        summed = 0
        for value in values:
            summed += value
    
        print("Worker %i: sum value is %i" % (uid, summed))
    
        return (uid, summed)
    
    
    if __name__ == '__main__':
        from random import sample
    
        # The process pool must be created inside __main__.
        async.pool = Pool(4)
    
        p = range(0, 1000)
        results = []
        for i in range(4):
            result = printsum(i, sample(p, 100))
            results.append(result)
    
        for result in results:
            print("Worker %i: sum value is %i" % result.get())
    

    在一个真实的案例中,我会详细说明装饰器,提供一些关闭它以进行调试的方法(同时保持未来的接口到位),或者可能是一种处理异常的工具;但我认为这很好地证明了这个原则。

    【讨论】:

    • 这应该是最好的答案。我喜欢它如何返回价值。不像只是异步运行的线程。
    【解决方案7】:

    只是

    import threading, time
    
    def f():
        print "f started"
        time.sleep(3)
        print "f finished"
    
    threading.Thread(target=f).start()
    

    【讨论】:

      【解决方案8】:

      你可以使用 eventlet。它可以让你编写看似同步的代码,但让它在网络上异步运行。

      这是一个超极小爬虫的例子:

      urls = ["http://www.google.com/intl/en_ALL/images/logo.gif",
           "https://wiki.secondlife.com/w/images/secondlife.jpg",
           "http://us.i1.yimg.com/us.yimg.com/i/ww/beta/y3.gif"]
      
      import eventlet
      from eventlet.green import urllib2
      
      def fetch(url):
      
        return urllib2.urlopen(url).read()
      
      pool = eventlet.GreenPool()
      
      for body in pool.imap(fetch, urls):
        print "got body", len(body)
      

      【讨论】:

        【解决方案9】:

        这样的东西对我有用,然后你可以调用该函数,它会将自己分派到一个新线程上。

        from thread import start_new_thread
        
        def dowork(asynchronous=True):
            if asynchronous:
                args = (False)
                start_new_thread(dowork,args) #Call itself on a new thread.
            else:
                while True:
                    #do something...
                    time.sleep(60) #sleep for a minute
            return
        

        【讨论】:

          【解决方案10】:

          从 Python 3.5 开始,您可以将增强的生成器用于异步函数。

          import asyncio
          import datetime
          

          增强的生成器语法:

          @asyncio.coroutine
          def display_date(loop):
              end_time = loop.time() + 5.0
              while True:
                  print(datetime.datetime.now())
                  if (loop.time() + 1.0) >= end_time:
                      break
                  yield from asyncio.sleep(1)
          
          
          loop = asyncio.get_event_loop()
          # Blocking call which returns when the display_date() coroutine is done
          loop.run_until_complete(display_date(loop))
          loop.close()
          

          新的async/await 语法:

          async def display_date(loop):
              end_time = loop.time() + 5.0
              while True:
                  print(datetime.datetime.now())
                  if (loop.time() + 1.0) >= end_time:
                      break
                  await asyncio.sleep(1)
          
          
          loop = asyncio.get_event_loop()
          # Blocking call which returns when the display_date() coroutine is done
          loop.run_until_complete(display_date(loop))
          loop.close()
          

          【讨论】:

          • @carnabeh,您能否扩展该示例以包含 OP 的“def longComputation()”函数?大多数示例使用“await asyncio.sleep(1)”,但如果 longComputation() 返回,例如,双精度,则不能只使用“await longComputation()”。
          • 未来十年,这应该是现在公认的答案。当你在 python3.5+ 中谈论 async 时,想到的应该是 asyncio 和 async 关键字。
          • 这个答案使用“新的和闪亮的”python 语法。现在这应该是排名第一的答案。
          【解决方案11】:

          您可以使用concurrent.futures(在 Python 3.2 中添加)。

          import time
          from concurrent.futures import ThreadPoolExecutor
          
          
          def long_computation(duration):
              for x in range(0, duration):
                  print(x)
                  time.sleep(1)
              return duration * 2
          
          
          print('Use polling')
          with ThreadPoolExecutor(max_workers=1) as executor:
              future = executor.submit(long_computation, 5)
              while not future.done():
                  print('waiting...')
                  time.sleep(0.5)
          
              print(future.result())
          
          print('Use callback')
          executor = ThreadPoolExecutor(max_workers=1)
          future = executor.submit(long_computation, 5)
          future.add_done_callback(lambda f: print(f.result()))
          
          print('waiting for callback')
          
          executor.shutdown(False)  # non-blocking
          
          print('shutdown invoked')
          

          【讨论】:

          • 这是一个非常棒的答案,因为它是这里唯一一个提供带有回调的线程池的可能性
          • 不幸的是,这也受到“全局解释器锁定”的影响。请参阅图书馆文档:link。使用 Python 3.7 测试
          • 这是一个阻塞异步调用
          【解决方案12】:

          您可以使用进程。如果你想在你的函数中永远运行它(比如网络):

          from multiprocessing import Process
          def foo():
              while 1:
                  # Do something
          
          p = Process(target = foo)
          p.start()
          

          如果您只想运行一次,请这样做:

          from multiprocessing import Process
          def foo():
              # Do something
          
          p = Process(target = foo)
          p.start()
          p.join()
          

          【讨论】:

            【解决方案13】:

            2021 年使用 Python 3.9 进行异步调用的原生 Python 方式也适用于 Jupyter / Ipython Kernel

            Camabeh 的答案是自 Python 3.3 以来要走的路。

            async def display_date(loop):
                end_time = loop.time() + 5.0
                while True:
                    print(datetime.datetime.now())
                    if (loop.time() + 1.0) >= end_time:
                        break
                    await asyncio.sleep(1)
            
            
            loop = asyncio.get_event_loop()
            # Blocking call which returns when the display_date() coroutine is done
            loop.run_until_complete(display_date(loop))
            loop.close()
            

            这将在 Jupyter Notebook / Jupyter Lab 中工作,但会引发错误:

            RuntimeError: This event loop is already running
            

            由于 Ipython 使用事件循环,我们需要一种称为嵌套异步循环的东西,它不是 yet implemented in Python。幸运的是有nest_asyncio 来处理这个问题。您需要做的就是:

            !pip install nest_asyncio # use ! within Jupyter Notebook, else pip install in shell
            import nest_asyncio
            nest_asyncio.apply()
            

            (基于this thread)

            只有当你调用 loop.close() 时它才会抛出另一个错误,因为它可能是指 Ipython 的主循环。

            RuntimeError: Cannot close a running event loop
            

            只要有人回复this github issue,我就会更新这个答案。

            【讨论】:

              猜你喜欢
              • 1970-01-01
              • 1970-01-01
              • 1970-01-01
              • 1970-01-01
              • 2020-12-19
              • 1970-01-01
              • 1970-01-01
              • 2015-11-28
              相关资源
              最近更新 更多