为何处理15k条记录时使用Threading仍卡顿?求优化方案
优化15k条链接重定向获取脚本的方案
我编写了一个Python脚本,功能为从CSV文件读取链接,访问后获取重定向链接并回存。但处理15k条记录时,即使已使用Threading,运行速度仍极慢。尝试多种方法后仍无改善,恳请专家提供优化方案。
原脚本代码
import concurrent.futures import sys import pandas as pd import requests from threading import Thread from queue import Queue out_put_file="" linkes = None out = [] urls = [] old = [] file_name =None concurrent = 10000 q = None count=0 df =None def do_work(): while True: global q url = q.get() res = get_status(url) q.task_done() def get_status(o_url): try: res = requests.get(o_url) if res: out.append(res.url) old.append(o_url) print(count) count=count+1 return [res.status_code,res.url ,o_url] except: pass return [ans.status_code,ans.url,url] def process_data(): global q global file_name global linkes global df file_name = input("Enter file name : ") file_name = file_name.strip() print("Generating .......") df = pd.read_csv(file_name+".csv") old_links =df["shopify"] for i in old_links: if type(i)!=str: urls.append(i) continue if not i.startswith("http"): linkes = "http://"+i urls.append(linkes) else: urls.append(i) df["shopify"]=urls q = Queue(concurrent * 2) for i in range(concurrent): t = Thread(target=do_work) t.daemon = True t.start() try: for url in urls: if type(url)!=str: continue q.put(url.strip()) q.join() except KeyboardInterrupt: sys.exit(1) process_data() for i in range (len(df['shopify'])): for j in range(len(old)): if df['shopify'][i]==old[j]: df['shopify'][i]=out[j] df = df[~df['shopify'].astype(str).str.startswith('http:')] df = df.dropna() df.to_csv(file_name+"-new.csv",index=False)
示例CSV数据
Email,shopify,Proofy_Status_Name hello@knobblystudio.com,http://puravidabracelets.myshopify.com,Deliverable service@cafe-select.co.uk,cafe-select.co.uk,Deliverable mtafich@gmail.com,,Deliverable whoopies@stevessnacks.com,stevessnacks.com,Deliverable customerservice@runwayriches.com,runwayriches.com,Deliverable shop@blackdogride.com.au,blackdogride.com.au,Deliverable anavasconcelos.nica@gmail.com,grass4you.com,Deliverable info@prideandprestigehair.com,prideandprestigehair.com,Deliverable info@dancinwoofs.com,dancinwoofs.com,Deliverable
问题分析与优化方案
原脚本核心问题
- 线程数设置过高:设为10000会导致系统上下文切换开销剧增,远超网络请求耗时,反而拖慢速度;同时requests的连接池无法支撑如此多并发连接。
- 线程不安全的全局变量:
out、old、count无锁保护,多线程同时操作会导致数据混乱、计数错误。 - 低效的数据匹配:最后用两层循环匹配新旧链接,O(n²)复杂度,15k条数据会产生2亿多次循环,耗时极长。
- 无连接复用:每次请求新建TCP连接,握手开销占比大。
- 异常处理不规范:
except:捕获所有异常,包括中断信号;出错后返回未定义的ans变量,会引发报错。 - 不必要的响应体下载:默认会下载完整响应体,而我们只需要重定向后的URL,浪费带宽和时间。
优化后的实现方案
import sys import pandas as pd import requests from concurrent.futures import ThreadPoolExecutor, as_completed def init_session(): # 创建会话,复用连接,减少TCP握手开销 session = requests.Session() # 设置请求头,模拟浏览器,避免被拦截 session.headers.update({ "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/114.0.0.0 Safari/537.36" }) return session def get_redirect_url(session, original_url): if not isinstance(original_url, str): return original_url, original_url # 处理无协议的URL,优先尝试HTTPS url = original_url if original_url.startswith(("http://", "https://")) else f"https://{original_url}" try: response = session.get(url, allow_redirects=True, stream=True, timeout=10) # 关闭连接,避免资源泄漏 response.close() return original_url, response.url except requests.exceptions.RequestException as e: # 若HTTPS失败,尝试HTTP if url.startswith("https://"): try: http_url = url.replace("https://", "http://") response = session.get(http_url, allow_redirects=True, stream=True, timeout=10) response.close() return original_url, response.url except requests.exceptions.RequestException: return original_url, original_url return original_url, original_url def main(): file_name = input("Enter file name : ").strip() print("Processing...") # 读取CSV df = pd.read_csv(f"{file_name}.csv") # 提取需要处理的链接 original_links = df["shopify"].tolist() # 初始化会话池 session = init_session() # 设置合理线程数(建议20-50,根据网络带宽调整) max_workers = 30 # 存储结果映射:原URL -> 重定向后的URL redirect_map = {} # 使用线程池处理请求 with ThreadPoolExecutor(max_workers=max_workers) as executor: # 提交所有任务 futures = {executor.submit(get_redirect_url, session, url): url for url in original_links} # 处理完成的任务 for idx, future in enumerate(as_completed(futures), 1): original_url, redirect_url = future.result() redirect_map[original_url] = redirect_url # 打印进度 if idx % 100 == 0: print(f"Processed {idx}/{len(original_links)} links") # 批量替换链接(O(n)复杂度) df["shopify"] = df["shopify"].map(lambda x: redirect_map.get(x, x)) # 过滤处理后的结果 df = df[~df["shopify"].astype(str).str.startswith("http:")] df = df.dropna() # 保存结果 df.to_csv(f"{file_name}-new.csv", index=False) print(f"Done! Results saved to {file_name}-new.csv") if __name__ == "__main__": try: main() except KeyboardInterrupt: sys.exit(1)
关键优化点说明
- 合理线程池配置:使用
ThreadPoolExecutor管理线程,设置20-50的线程数,平衡并发与系统开销。 - 连接复用:用
requests.Session()复用TCP连接,大幅减少握手耗时。 - 线程安全的结果存储:通过任务返回结果构建字典映射,避免全局变量的线程安全问题。
- 高效数据替换:用字典映射批量替换链接,O(n)复杂度,替代原有的O(n²)循环。
- 减少不必要的下载:设置
stream=True只获取响应头,不下载响应体,节省带宽和时间。 - 完善的异常处理:只捕获requests相关异常,同时尝试HTTPS/HTTP降级,提升成功率。
- 进度反馈:每处理100条打印进度,方便监控状态。
- 代码结构优化:去掉冗余全局变量,使用函数封装,逻辑更清晰。
内容的提问来源于stack exchange,提问作者Nabi Bux
相关产品推荐
相关产品推荐

