Python多线程脚本挂起原因排查:Lock/Semaphore使用是否有误?
Python多线程脚本挂起问题排查与修复
问题现象
脚本在栈帧的wait_for_tstate_lock elif lock.acquire(block, timeout)行持续挂起,怀疑Lock()或Semaphore()使用错误,相关代码及UrlFeaturizer实现如下:
原始主脚本
import concurrent.futures import csv import threading import pandas as pd from tqdm import tqdm from discovery.featurizer import UrlFeaturizer semaphore = threading.Semaphore(50) def run_featurizer(): res = UrlFeaturizer(url).run(dataset)[1] return res if __name__ == "__main__": keys = UrlFeaturizer("1.1.1.1").run("")[0] with open("num_features.csv", "a") as f: csv_out = csv.DictWriter(f, keys) csv_out.writeheader() with semaphore: with concurrent.futures.ThreadPoolExecutor(max_workers=50) as executor: futures = [] datasets = [ "benign_domains.csv", "dmca_domains.csv", ] for dataset in datasets: urls = pd.read_csv(dataset, header=None).iloc[:, 0].to_list() with tqdm(total=len(urls)) as pbar: for url in urls: futures.append(executor.submit(run_featurizer, )) for future in concurrent.futures.as_completed(futures): pbar.update() with threading.Lock() as lock: csv_out.writerow(future.result()) f.flush()
原始UrlFeaturizer实现
class UrlFeaturizer(object): def __init__(self, url): self.url = url try: self.response = requests.get( prepend_protocols(self.url), headers=headers, timeout=5 ) except Exception: self.response = None try: self.whois = whois.query(self.url).__dict__ except Exception: self.whois = None try: self.soup_c = BeautifulSoup( self.response.content, features="lxml", from_encoding=self.response.encoding, ) except Exception: self.soup_c = None def lookup_whois(self) -> int: return int(False) if self.whois else int(True) def lookup_domain_age(self) -> int: if self.whois and self.whois["creation_date"]: return (date.today() - self.whois["creation_date"].date()).days return def verify_ssl(self) -> bool: try: ssl_cert = ssl.get_server_certificate((self.url, 443), timeout=10) return int(True) if ssl_cert else int(False) except Exception: return def check_security(self) -> bool: try: requests.head(f"https://{self.url}", timeout=10) return int(True) except Exception: return int(False) def has_com_tld(self): return int(True) if extract_tld(self.url) == "com" else int(False) def run(self, dataset=None): data = { "url": self.url, "uses_whois_privacy": self.lookup_whois(), "domain_age": self.lookup_domain_age(), "has_ssl": self.verify_ssl(), "is_secure": self.check_security(), "has_com_tld": self.has_com_tld(), "label": Path(dataset).stem, } return data.keys(), data
核心问题分析
- 全局变量引用混乱:
run_featurizer直接引用外部循环的url和dataset变量,多线程下会导致变量值被覆盖,引发不可预期的错误。 - Semaphore完全误用:将
semaphore的with语句放在线程池外层,相当于只允许一个线程池实例运行,彻底失去并发控制作用,甚至可能触发死锁。 - 无效锁机制:每次处理结果时创建新的
threading.Lock(),每个线程持有独立的锁,完全无法保护共享的文件写入操作,反而增加不必要的开销。 - 进度条逻辑错误:在每个数据集循环内初始化
tqdm,但futures列表是累积的,会导致进度条更新次数远超实际URL数量。 - UrlFeaturizer阻塞风险:
__init__中直接执行多个网络请求,异常处理过于宽泛,单个请求超时或失败可能导致线程长时间阻塞,拖垮整个线程池。
修复后的代码
主脚本
import concurrent.futures import csv import threading from pathlib import Path import pandas as pd from tqdm import tqdm from discovery.featurizer import UrlFeaturizer # 全局锁:保护共享文件写入操作 write_lock = threading.Lock() def run_featurizer(url, dataset): # 通过参数传递变量,避免全局引用冲突 res = UrlFeaturizer(url).run(dataset)[1] return res if __name__ == "__main__": keys = UrlFeaturizer("1.1.1.1").run("")[0] # 使用w模式避免重复写入表头,添加newline防止CSV空行 with open("num_features.csv", "w", newline="") as f: csv_out = csv.DictWriter(f, keys) csv_out.writeheader() # 线程池max_workers直接控制并发,无需额外Semaphore with concurrent.futures.ThreadPoolExecutor(max_workers=50) as executor: futures = [] datasets = ["benign_domains.csv", "dmca_domains.csv"] total_urls = 0 # 先统计总URL数,用于全局进度条 for dataset in datasets: urls = pd.read_csv(dataset, header=None).iloc[:, 0].to_list() total_urls += len(urls) for url in urls: futures.append(executor.submit(run_featurizer, url, dataset)) # 全局进度条,确保更新准确 with tqdm(total=total_urls) as pbar: for future in concurrent.futures.as_completed(futures): pbar.update(1) with write_lock: try: result = future.result() csv_out.writerow(result) f.flush() except Exception as e: # 捕获任务异常,避免程序崩溃 print(f"处理URL失败: {e}")
优化后的UrlFeaturizer
import ssl import requests from bs4 import BeautifulSoup import whois from datetime import date from tldextract import extract_tld def prepend_protocols(url): if not url.startswith(("http://", "https://")): return f"https://{url}" return url headers = {"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/91.0.4472.124 Safari/537.36"} class UrlFeaturizer(object): def __init__(self, url): self.url = url self.response = None self.whois = None self.soup_c = None # 拆分网络操作,避免__init__阻塞过久 self._fetch_response() self._fetch_whois() self._parse_html() def _fetch_response(self): try: self.response = requests.get( prepend_protocols(self.url), headers=headers, timeout=5 ) self.response.raise_for_status() except requests.exceptions.RequestException as e: print(f"请求URL失败: {self.url}, 错误: {e}") def _fetch_whois(self): try: whois_result = whois.query(self.url) if whois_result: self.whois = whois_result.__dict__ except Exception as e: print(f"WHOIS查询失败: {self.url}, 错误: {e}") def _parse_html(self): if self.response and self.response.content: try: self.soup_c = BeautifulSoup( self.response.content, features="lxml", from_encoding=self.response.encoding ) except Exception as e: print(f"解析HTML失败: {self.url}, 错误: {e}") def lookup_whois(self) -> int: return int(False) if self.whois else int(True) def lookup_domain_age(self) -> int: if self.whois and self.whois.get("creation_date"): creation_date = self.whois["creation_date"] # 处理whois返回日期列表的情况 if isinstance(creation_date, list): creation_date = creation_date[0] return (date.today() - creation_date.date()).days return 0 # 返回默认值,避免None导致CSV写入错误 def verify_ssl(self) -> int: try: ssl_cert = ssl.get_server_certificate((self.url, 443), timeout=10) return int(True) if ssl_cert else int(False) except Exception as e: print(f"SSL验证失败: {self.url}, 错误: {e}") return int(False) def check_security(self) -> int: try: response = requests.head(f"https://{self.url}", timeout=10, allow_redirects=True) response.raise_for_status() return int(True) except Exception as e: print(f"安全检查失败: {self.url}, 错误: {e}") return int(False) def has_com_tld(self) -> int: try: tld = extract_tld(self.url).suffix return int(True) if tld == "com" else int(False) except Exception as e: print(f"提取TLD失败: {self.url}, 错误: {e}") return int(False) def run(self, dataset=None): data = { "url": self.url, "uses_whois_privacy": self.lookup_whois(), "domain_age": self.lookup_domain_age(), "has_ssl": self.verify_ssl(), "is_secure": self.check_security(), "has_com_tld": self.has_com_tld(), "label": Path(dataset).stem if dataset else "" } return data.keys(), data
关键修复说明
- 移除无效Semaphore:
ThreadPoolExecutor的max_workers已实现并发控制,原Semaphore的使用方式完全错误,直接移除。 - 参数化传递变量:
run_featurizer通过参数接收url和dataset,彻底解决多线程下的变量引用冲突。 - 全局锁保护写入:使用单个全局
threading.Lock()确保同一时间只有一个线程写入文件,避免数据错乱。 - 修正进度条逻辑:统计所有URL总数后创建全局进度条,确保进度更新准确。
- 优化UrlFeaturizer:拆分网络操作、增强异常捕获、设置返回值默认值,降低线程阻塞风险,提升鲁棒性。
- 异常处理增强:捕获任务执行异常,避免单个URL处理失败导致整个程序崩溃。
内容的提问来源于stack exchange,提问作者ariyasas94
相关产品推荐
相关产品推荐

