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

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

核心问题分析

  1. 全局变量引用混乱:run_featurizer直接引用外部循环的url和dataset变量,多线程下会导致变量值被覆盖,引发不可预期的错误。
  2. Semaphore完全误用:将semaphore的with语句放在线程池外层,相当于只允许一个线程池实例运行,彻底失去并发控制作用,甚至可能触发死锁。
  3. 无效锁机制:每次处理结果时创建新的threading.Lock(),每个线程持有独立的锁,完全无法保护共享的文件写入操作,反而增加不必要的开销。
  4. 进度条逻辑错误:在每个数据集循环内初始化tqdm,但futures列表是累积的,会导致进度条更新次数远超实际URL数量。
  5. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 16:00:59