使用concurrent.futures和/或asyncio进行大规模并发处理的方案、问题及优化问询
使用concurrent.futures和/或asyncio进行大规模并发处理的方案、问题及优化问询
我现在要在数千线程的线程池上运行数十亿个IO密集型任务,目前碰到了这些核心挑战:
- 降低内存占用:
concurrent.futures的线程池用的是无界队列,提交任务太多会导致内存暴涨,但提交太少又完不成足够的工作量。 - 收集结果:不管是线程池还是它的上下文管理器都不会存储已完成的Future,不过收集结果其实不是降低内存占用的必要操作,但我希望之后如果需要的话能有收集结果的选项。
- 异常处理:我希望任务里抛出的Python异常能传递到主线程。
- 整合asyncio:要通过
loop.run_in_executor()把asyncio结合进来。
关于concurrent.futures的初步处理思路:每个Future都可以通过add_done_callback来收集结果并处理异常。不过回调里抛出的异常会被忽略,但好在回调是在主线程运行的,所以可以把异常存到全局变量里,然后在提交任务的循环里重新抛出。不过这对降低内存占用没帮助,而且就算用asyncio我也还没搞懂怎么解决内存问题(后面会说到)。
关于整合asyncio的尝试:我创建了一个消费者任务和一个生产者任务。生产者在同步循环里把可调用对象提交给事件循环的执行器,然后把得到的Future放到一个无界队列里。
消费者会维护一个待处理Future的集合,等待这个集合或者队列里的第一个结果。它会同时等待两个对象的结果:一个是队列,另一个是对这个待处理Future集合的asyncio.wait()协程。如果结果来自队列,就把这个Future加到集合里继续等待;如果是已完成的Future,就先处理结果、处理异常,然后再继续等待。等队列关闭后,这个任务就完成了。
我是不是把这件事搞得太复杂了?有没有更简单的实现方式?
另外我还没找到一个优雅的降低内存占用的方法。目前能想到的最好办法就是检查线程池的_work_queue.qsize(),如果超过阈值就await asyncio.sleep(<一小段时间>)。肯定有更好的办法吧?
备注:内容来源于stack exchange,提问作者user19087
相关产品推荐
相关产品推荐

