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)。
在实施您的解决方案时,您可能会发现您想要修改现有代码以使用协同程序而不是线程。