Python3中concurrent.futures.ThreadPoolExecutor的内存使用情况

gag*_*env 6 python malloc concurrency json urllib

我正在构建一个脚本来下载和解析奥巴马医疗保健交易所健康保险计划的福利信息.部分原因是需要从每个保险公司下载和解析计划权益JSON文件.为了做到这一点,我使用concurrent.futures.ThreadPoolExecutor6个工作人员下载每个文件(使用urllib),解析并循环通过JSON并提取相关信息(存储在脚本中的嵌套字典中).

(在win32上运行Python 3.5.1(v3.5.1:37a07cee5969,2015年12月6日,01:38:48)[MSC v.1900 32位(英特尔)]

问题是,当我同时执行此操作时,脚本在通过JSON文件下载\ parsed\loop后似乎不会释放内存,过了一会儿,它会崩溃,malloc引发内存错误.

当我连续地进行 - 用一个简单的for in循环 - 然而,程序不会崩溃,也不会占用极大的内存.

def load_json_url(url, timeout):
    req = urllib.request.Request(url, headers={ 'User-Agent' : 'Mozilla/5.0' })
    resp = urllib.request.urlopen(req).read().decode('utf8')
    return json.loads(resp) 



 with concurrent.futures.ThreadPoolExecutor(max_workers=6) as executor:
        # Start the load operations and mark each future with its URL
        future_to_url = {executor.submit(load_json_url, url, 60): url for url in formulary_urls}
        for future in concurrent.futures.as_completed(future_to_url):
            url = future_to_url[future]
            try:
                # The below timeout isn't raising the TimeoutError.
                data = future.result(timeout=0.01)
                for item in data:
                        if item['rxnorm_id']==drugid: 
                            for row in item['plans']:
                                print (row['drug_tier'])
                                (plansid_dict[row['plan_id']])['drug_tier']=row['drug_tier']
                                (plansid_dict[row['plan_id']])['prior_authorization']=row['prior_authorization']
                                (plansid_dict[row['plan_id']])['step_therapy']=row['step_therapy']
                                (plansid_dict[row['plan_id']])['quantity_limit']=row['quantity_limit']

            except Exception as exc:
                print('%r generated an exception: %s' % (url, exc))


            else:
                downloaded_plans=downloaded_plans+1
Run Code Online (Sandbox Code Playgroud)

pre*_*per 7

作为替代解决方案,您可以调用add_done_callback期货而根本不使用as_completed。关键是不要保留对期货的引用。所以future_to_url原始问题中的列表是一个坏主意。

我所做的基本上是:

def do_stuff(future):
    res = future.result()  # handle exceptions here if you need to

f = executor.submit(...)
f.add_done_callback(do_stuff)
Run Code Online (Sandbox Code Playgroud)


Mow*_*hon 6

如果您使用标准模块 \xe2\x80\x9cconcurrent.futures\xe2\x80\x9d 并希望同时处理数百万数据,那么一个工作队列将占用所有可用内存。

\n\n

您可以使用bounded-pool-executor。\n https://github.com/mowshon/bounded_pool_executor

\n\n
pip install bounded-pool-executor\n
Run Code Online (Sandbox Code Playgroud)\n\n

例子:

\n\n
from bounded_pool_executor import BoundedProcessPoolExecutor\nfrom time import sleep\nfrom random import randint\n\ndef do_job(num):\n    sleep_sec = randint(1, 10)\n    print(\'value: %d, sleep: %d sec.\' % (num, sleep_sec))\n    sleep(sleep_sec)\n\nwith BoundedProcessPoolExecutor(max_workers=5) as worker:\n    for num in range(10000):\n        print(\'#%d Worker initialization\' % num)\n        worker.submit(do_job, num)\n
Run Code Online (Sandbox Code Playgroud)\n


小智 5

这不是你的错。as_complete()在完成之前不会释放其期货。已经记录了一个问题:https : //bugs.python.org/issue27144

就目前而言,我认为多数方法是将as_complete()封装在另一个块中,将其分块为同等数量的期货,具体取决于您要花费多少RAM和结果多少。它会阻塞在每个块上,直到所有工作都消失,然后再进入下一个块为止,这样会变慢或长时间停留在中间,但是我现在看不到其他任何方法,尽管在有更聪明的时候会保留此答案道路。