如何用Dask并行化Python+Selenium爬取巴西SISAB健康数据
用Dask并行优化SISAB数据爬取方案
核心选择:用dask.delayed而非Dask DataFrame
Dask DataFrame更适合处理已有的结构化数据分片,而你的场景是并行执行独立的爬取任务,每个任务参数独立、无依赖,dask.delayed是更灵活的选择——它能直接包装你的现有Selenium爬取函数,自动实现并行调度。
具体实现步骤
1. 准备参数组合列表
先把所有17000种参数组合整理成Python列表,每个元素是包含诊疗代码、性别、年龄组、年月的元组或字典,示例:
params_list = [ {"code": "CID-XXX", "gender": "男", "age_group": "0-4岁", "year_month": "2018-01"}, # ... 其他所有参数组合 ]
2. 封装独立爬取函数
把你现有的串行爬取逻辑封装成独立函数,输入一组参数,返回爬取到的DataFrame(或直接保存临时文件):
from selenium import webdriver from selenium.webdriver.support.ui import WebDriverWait from selenium.webdriver.support import expected_conditions as EC import pandas as pd def scrape_single_param(param): # 每个任务独立初始化浏览器,避免实例共享冲突 options = webdriver.ChromeOptions() options.add_argument("--headless=new") driver = webdriver.Chrome(options=options) try: driver.get("https://sisab.saude.gov.br/paginas/acessoRestrito/relatorio/federal/saude/RelSauProducao.xhtml") # 设置筛选条件(替换成你实际的元素定位逻辑) driver.find_element("id", "codigoInput").send_keys(param["code"]) driver.find_element("id", "generoSelect").select_by_visible_text(param["gender"]) driver.find_element("id", "faixaEtariaSelect").select_by_visible_text(param["age_group"]) driver.find_element("id", "mesAnoInput").send_keys(param["year_month"]) # 提交查询并等待结果加载 driver.find_element("id", "consultarBtn").click() WebDriverWait(driver, 10).until(EC.presence_of_element_located(("id", "resultadoTabela"))) # 提取表格数据转成DataFrame,添加参数标识列 table = driver.find_element("id", "resultadoTabela") df = pd.read_html(table.get_attribute("outerHTML"))[0] df["诊疗代码"] = param["code"] df["性别"] = param["gender"] df["年龄组"] = param["age_group"] df["年月"] = param["year_month"] return df except Exception as e: # 异常捕获,避免单个任务失败中断全局流程 print(f"参数{param}爬取失败: {str(e)}") return pd.DataFrame() finally: driver.quit()
3. 用dask.delayed包装任务
将每个参数对应的爬取逻辑包装为延迟执行对象:
from dask import delayed delayed_tasks = [delayed(scrape_single_param)(param) for param in params_list]
4. 启动集群并行执行
根据机器资源设置并行进程数(注意不要触发目标网站反爬),执行任务并合并结果:
from dask.distributed import Client import dask # 启动本地集群,建议进程数设为CPU核心数的1-2倍(避免反爬或资源过载) client = Client(n_workers=6, threads_per_worker=1) # 执行所有延迟任务,得到结果列表 results = dask.compute(*delayed_tasks) # 合并所有有效结果 final_df = pd.concat([df for df in results if not df.empty], ignore_index=True) # 用Parquet格式保存(比Excel更高效,适合大数据量) final_df.to_parquet("sisab_anxiety_depression_data.parquet", index=False) client.close()
关键注意事项
- 反爬规避:并行爬取会提升请求频率,建议给每个任务添加1-3秒的随机延迟,限制并行数从低到高测试,避免高峰时段爬取。
- 资源控制:每个Selenium实例会占用内存,并行数不要超过机器承载上限,否则会导致卡顿或崩溃。
- 容错机制:给爬取函数添加重试逻辑(比如失败后重试1-2次),避免单个任务失败影响全局。
- 内存优化:如果内存不足以存储所有结果,可以让每个爬取任务直接保存为单独的Parquet/CSV文件,最后再批量合并,避免内存溢出。
内容的提问来源于stack exchange,提问作者KVemuri
相关产品推荐
相关产品推荐

