Python使用concurrent.futures做多线程爬虫报ValueError列不匹配问题
问题原因
- 表头提取逻辑错误:
process_data中遍历所有<tr>行时,cols变量会被后续无表头的行覆盖,只有第一行<tr>有<th>表头元素,后续行的row.find_all('th')返回空列表,最终column[0]是空列表,执行column[0].append('Currency')后,columns参数只有1个元素,和数据的7列不匹配,直接触发报错。
- 表头提取逻辑错误:
- 数据收集逻辑错误:单线程版本会收集每个页面所有
<tr>对应的行数据,而改造后的process_data中outs变量在循环中被反复覆盖,最终每个链接只返回最后一行的数据,导致output列表长度远小于预期。
- 数据收集逻辑错误:单线程版本会收集每个页面所有
- 并发设计冗余且存在隐患:IO密集型的爬虫/网页解析任务无需同时用进程池+线程池,且
requests.Session对象无法跨进程序列化传递,进程池的使用完全多余,还会带来额外的序列化开销和不可预期的错误。
- 并发设计冗余且存在隐患:IO密集型的爬虫/网页解析任务无需同时用进程池+线程池,且
- 冗余删除逻辑:单线程中每个页面处理完删output第一条(表头对应的空td行),多线程版本只删了一次output第一条,就算前面逻辑对也会有数据错位问题。
修正后代码
import requests from bs4 import BeautifulSoup import pandas as pd import time import concurrent.futures def process_data(key, page_content): soup = BeautifulSoup(page_content, 'html.parser') table = soup.select('table')[0] rows = table.find_all('tr') # 仅提取一次表头,存入函数属性避免重复提取 if not hasattr(process_data, 'columns'): cols = [item.text.strip() for item in rows[0].find_all('th')] cols.append('Currency') process_data.columns = cols # 收集当前页面所有有效数据行,跳过表头行 page_data = [] for row in rows[1:]: outs = [item.text.strip() for item in row.find_all('td')] if outs: outs.append(key) page_data.append(outs) return page_data def getCurrencyHistorical(session, item): key, url = item resp = session.get(url) resp.raise_for_status() return process_data(key, resp.content) def main(): t1 = time.perf_counter() links = { "USD-IDR":"https://www.investing.com/currencies/usd-idr-historical-data", "USD-JPY":"https://www.investing.com/currencies/usd-jpy-historical-data", "USD-CNY":"https://www.investing.com/currencies/usd-cny-historical-data" } # 补全请求头避免被网站反爬拦截 headers = { 'Accept-Language': 'en-US,en;q=0.9', 'Upgrade-Insecure-Requests': '1', 'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/88.0.4324.150 Safari/537.36 Edg/88.0.705.63', 'Accept': 'text/html,application/xhtml+xml,application/xml;q=0.9,image/webp,image/apng,*/*;q=0.8,application/signed-exchange;v=b3;q=0.9', 'Cache-Control': 'max-age=0', 'Connection': 'keep-alive' } output = [] with requests.Session() as session: session.headers.update(headers) # IO密集型任务仅用线程池即可,无需额外进程池开销 with concurrent.futures.ThreadPoolExecutor(max_workers=len(links)) as executor: for page_data in executor.map(lambda x: getCurrencyHistorical(session, x), links.items()): output.extend(page_data) df = pd.DataFrame(output, columns=process_data.columns) t2 = time.perf_counter() print(f'Finished in {t2-t1:.2f} seconds') print(df) return df if __name__ == '__main__': main()
内容的提问来源于stack exchange,提问作者ohai
相关产品推荐
相关产品推荐

