You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Python爬虫中使用concurrent.futures实现并发爬取并优化现有代码

改造思路
  • 爬虫属于典型IO密集型任务,优先选用concurrent.futures.ThreadPoolExecutor(线程池)实现并发,相比多进程模式开销更低,适配网络请求等待场景
  • 将单个国家的爬取、初步解析逻辑抽离为独立可并发执行的函数
  • 单个任务处理时直接为返回数据新增Country列,避免并发返回顺序和提交顺序不一致导致的国家数据匹配错误
  • 公共的请求头、Cookie、参数模板、列名处理逻辑保留在主函数中,避免重复计算
改造后完整代码
import json
import pandas as pd
import copy
import requests
from bs4 import BeautifulSoup
import concurrent.futures

# 抽离单个国家的爬取逻辑,作为线程池提交的独立任务
def fetch_single_country(params, country_name, headers, cookies):
    try:
        response = requests.post(
            'https://satudata.kemendag.go.id/balance-query',
            headers=headers,
            cookies=cookies,
            data=params,
            timeout=30
        )
        response.raise_for_status()
        json_data = json.loads(response.text)
        df1 = pd.json_normalize(json_data)
        table_list = df1['data']
        country_df = pd.DataFrame()
        for table in table_list:
            df = pd.DataFrame.from_records(table)
            country_df = pd.concat([country_df, df], ignore_index=True)
        # 直接给当前国家数据加Country列,避免后续匹配错误
        country_df['Country'] = country_name
        return country_df
    except Exception as e:
        print(f"爬取{country_name}失败,错误信息:{e}")
        return pd.DataFrame()

def getBalanceofTrade():
    cookies = {
        '_ga': 'GA1.3.672481365.1630462756',
        'sc_is_visitor_unique': 'rx12421317.1631602001.FE09E031FA464FEEAD095B438DD92829.2.2.2.2.2.2.2.2.2',
        '_ga_41205092': 'GS1.1.1631601995.2.0.1631602551.0',
        'varient_csrf_cookie': 'c34f577f894740d92106453d293f69b4',
        'ci_session': 'q9qvice0p2dgnp5f50veghgrqeveavrm',
        '_gid': 'GA1.3.542540031.1635307774',
        '_gat_gtag_UA_155890962_1': '1',
        'f5avr0338662302aaaaaaaaaaaaaaaa_cspm_': 'DMPCJNDKGKPBLIDLCKDBLLDGAEPJPBKDHBPPMIDGKFDEMKGDHFNEBJBDOLAMKFJACFDCDGKEHKFEFIPPNNKAJEACAPLNMFADDKOEHBJBKJGGDMKBNHHFMJBOEGLPMEPA',
    }

    headers = {
        'Connection': 'keep-alive',
        'sec-ch-ua': '"Chromium";v="94", "Google Chrome";v="94", ";Not A Brand";v="99"',
        'Accept': 'application/json, text/javascript, */*; q=0.01',
        'Content-Type': 'application/x-www-form-urlencoded; charset=UTF-8',
        'X-Requested-With': 'XMLHttpRequest',
        'sec-ch-ua-mobile': '?0',
        'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/94.0.4606.81 Safari/537.36',
        'sec-ch-ua-platform': '"Windows"',
        'Origin': 'https://satudata.kemendag.go.id',
        'Sec-Fetch-Site': 'same-origin',
        'Sec-Fetch-Mode': 'cors',
        'Sec-Fetch-Dest': 'empty',
        'Referer': 'https://satudata.kemendag.go.id/balance-of-trade-with-trade-partner-country',
        'Accept-Language': 'en-US,en;q=0.9',
    }
    
    base_params = {
        'varient_csrf_token': 'c34f577f894740d92106453d293f69b4',
        'category': 'country'
    }
    
    value_id = pd.read_excel('Value ID Countries.xlsx')
    value_dict = dict(zip(value_id['VALUE ID'], value_id['COUNTRY']))
    
    # 线程池并发爬取所有国家数据
    df_all = pd.DataFrame()
    max_workers = 8 # 可根据实际情况调整,不要过大
    with concurrent.futures.ThreadPoolExecutor(max_workers=max_workers) as executor:
        future_list = []
        for region_id, country_name in value_dict.items():
            new_params = copy.copy(base_params)
            new_params['region_id'] = region_id
            future = executor.submit(fetch_single_country, new_params, country_name, headers, cookies)
            future_list.append(future)
        # 收集所有返回结果
        for future in concurrent.futures.as_completed(future_list):
            res_df = future.result()
            if not res_df.empty:
                df_all = pd.concat([df_all, res_df], ignore_index=True)
    
    # 原有数据清洗逻辑
    df_new = df_all.iloc[:, 0:6]
    df_new1 = df_all.iloc[:, 8:9]
    df_new2 = df_all.iloc[:, 10:11]
    result = pd.concat([df_new, df_new1, df_new2], axis=1)
    
    URL = 'https://satudata.kemendag.go.id/balance-of-trade-with-trade-partner-country'
    page = requests.get(URL)
    soup = BeautifulSoup(page.content, 'html.parser')
    table = soup.find('table', id="table-balance")
    rows = table.find_all('th')
    cols = [item.text.strip() for item in rows]
    # 列名处理
    column_name = [cols[0]]
    column_name1 = cols[1:6]
    column_name2 = [cols[10]]
    tahun = column_name1 + column_name2
    final_column = column_name + column_name1 + column_name2 + ['Country']
    result.columns = final_column
    
    desc_1 = 3*(result['Uraian'][0],)
    desc_2 = 3*(result['Uraian'][3],)
    desc_3 = 3*(result['Uraian'][6],)
    desc_4 = 3*(result['Uraian'][9],)
    x = desc_1 + desc_2 + desc_3 + desc_4
    y = list(x)
    z = y * len(value_id['COUNTRY'])
    result = result.replace(['TOTAL PERDAGANGAN', 'EKSPOR', 'IMPOR', 'NERACA PERDAGANGAN'], ['Total', 'Total', 'Total', 'Total']) 
    df2 = pd.DataFrame(z, columns=['description'])
    
    result = result.rename(columns = {'Uraian': 'tipe'}, inplace = False)
    table = pd.concat ([df2, result], axis=1)
    
    df_table = table.melt(id_vars=['description', 'tipe', 'Country'], value_vars=tahun, var_name='year')
    
    return df_table
注意事项
  • max_workers不要设置过大,建议控制在3~10之间,避免请求频率过高触发网站反爬封禁
  • 代码中使用的Cookie、CSRF Token存在有效期,若后续请求返回失败,可重新访问目标网站后从浏览器开发者工具的网络请求中复制替换对应值
  • 可根据需要在fetch_single_country函数中添加重试、异常捕获逻辑,提升爬取成功率

内容的提问来源于stack exchange,提问作者ohai

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.27 21:36:03