【问题标题】:Return from function if execution finished within timeout or make callback otherwise如果执行在超时内完成,则从函数返回,否则进行回调
【发布时间】:2017-03-05 05:13:47
【问题描述】:

我有一个 Python 3.5 项目,没有使用任何异步功能。我必须实现以下逻辑:

def should_return_in_3_sec(some_serious_job, arguments, finished_callback):
    # Start some_serious_job(*arguments) in a task
    # if it finishes within 3 sec:
    #    return result immediately
    # otherwise return None, but do not terminate task.
    # If the task finishes in 1 minute:
    #    call finished_callback(result)
    # else:
    #    call finished_callback(None)
    pass

函数should_return_in_3_sec() 应该保持同步,但由我来编写任何新的异步代码(包括some_serious_job())。

最优雅和最Pythonic的方式是什么?

【问题讨论】:

    标签: python python-3.x asynchronous async-await python-asyncio


    【解决方案1】:

    fork 一个线程来完成重要的工作,让它将其结果写入队列,然后终止。从该队列中读取您的主线程,超时为三秒。如果发生超时,则启动另一个线程并返回 None。让第二个线程从队列中读取超时一分钟;如果也超时,请调用 finished_callback(None);否则调用finished_callback(result)。

    我是这样画的:

    import threading, queue
    
    def should_return_in_3_sec(some_serious_job, arguments, finished_callback):
      result_queue = queue.Queue(1)
    
      def do_serious_job_and_deliver_result():
        result = some_serious_job(arguments)
        result_queue.put(result)
    
      threading.Thread(target=do_serious_job_and_deliver_result).start()
    
      try:
        result = result_queue.get(timeout=3)
      except queue.Empty:  # timeout?
    
        def expect_and_handle_late_result():
          try:
            result = result_queue.get(timeout=60)
          except queue.Empty:
            finished_callback(None)
          else:
            finished_callback(result)
    
        threading.Thread(target=expect_and_handle_late_result).start()
        return None
      else:
        return result
    

    【讨论】:

    • 不错!虽然我没有在以后的超时中找到终止 Serious_job ...我也想知道,在这个多线程示例中使用标准 queue 是否可以?
    • queue 模块适用于线程上下文,是的。它显式地同步其使用。当然,您在终止分离线程时会遇到一个小问题(取决于您的上下文是否真的会成为问题)。但是你给他们的规格不允许再次加入他们。
    【解决方案2】:

    threading 模块有一些简单的超时选项,例如参见Thread.join(timeout)

    如果您确实选择使用 asyncio,以下是满足您某些需求的部分解决方案:

    import asyncio
    
    import time
    
    
    async def late_response(task, flag, timeout, callback):
        done, pending = await asyncio.wait([task], timeout=timeout)
        callback(done.pop().result() if done else None)  # will raise an exception if some_serious_job failed
        flag[0] = True  # signal some_serious_job to stop
        return await task
    
    
    async def launch_job(loop, some_serious_job, arguments, finished_callback,
                         timeout_1=3, timeout_2=5):
        flag = [False]
        task = loop.run_in_executor(None, some_serious_job, flag, *arguments)
        done, pending = await asyncio.wait([task], timeout=timeout_1)
        if done:
            return done.pop().result()  # will raise an exception if some_serious_job failed
        asyncio.ensure_future(
            late_response(task, flag, timeout_2, finished_callback))
        return None
    
    
    def f(flag, n):
        for i in range(n):
            print("serious", i, flag)
            if flag[0]:
                return "CANCELLED"
            time.sleep(1)
        return "OK"
    
    
    def finished(result):
        print("FINISHED", result)
    
    
    loop = asyncio.get_event_loop()
    result = loop.run_until_complete(launch_job(loop, f, [1], finished))
    print("result:", result)
    loop.run_forever()
    

    这将在单独的线程中运行作业(使用 loop.set_executor(ProcessPoolExecutor()) 在进程中运行 CPU 密集型任务)。请记住,终止进程/线程是一种不好的做法 - 上面的代码使用一个非常简单的列表来指示线程停止(另请参阅 threading.Event / multiprocessing.Event)。

    在实施您的解决方案时,您可能会发现您想要修改现有代码以使用协同程序而不是线程。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2016-02-17
      • 1970-01-01
      • 2012-04-20
      • 1970-01-01
      • 2021-08-02
      • 2019-08-18
      • 2012-05-25
      相关资源
      最近更新 更多