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

HTTPX与ThreadPoolExecutor结合时遇RuntimeError: Event loop is closed问题求助

问题描述

我正在将WordPress死链检测脚本从requests迁移至HTTPX——原requests脚本无法绕过Cloudflare网站防护,尝试多种方法都没解决,同时想测试异步特性能否提升脚本运行速度。

该脚本功能是扫描整个WordPress网站的死链并存入CSV文件,原版本用ThreadPoolExecutor批量检测单篇文章中的所有链接。现在我想用Asyncio/HTTPX实现相同功能,写的代码如下,但运行时出现错误:

原代码
import httpx
import asyncio
import bs4
import sys
import csv
import concurrent.futures
from concurrent.futures import as_completed

headers = {
        'accept': 'text/html,application/xhtml+xml,application/xml;q=0.9,image/avif,image/webp,image/apng,*/*;q=0.8,application/signed-exchange;v=b3;q=0.7',
        'accept-language': 'en-US,en-GB;q=0.9,en;q=0.8',
        'cache-control': 'max-age=0',
        'dnt': '1',
        'priority': 'u=0, i',
        'sec-fetch-dest': 'document',
        'sec-fetch-mode': 'navigate',
        'upgrade-insecure-requests': '1',
        'User-Agent': 'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/16.0 Safari/605.1.15 Edg/126.0.0.0'
    }
    
async def threadWorkAsync(link, headers, client):
        print("\t Content Link:", link[0], end=" ")
        try:
            link_response= await client.head(link[0], headers=headers)
            print(link_response.status_code)
        except (
            httpx.HTTPStatusError,
            httpx.HTTPError,
            httpx.ConnectError,
            httpx.InvalidURL,
            httpx.ConnectTimeout,
            httpx.RequestError,
        ) as errh:
            print("{errh} in URL, ", link)    

def asyncThreadRunner(link, headers, client):
     # Create run loop for this thread and block until completion
     asyncio.set_event_loop_policy(asyncio.WindowsSelectorEventLoopPolicy())
     asyncio.run(threadWorkAsync(link, headers, client))

def executeBrokenLinkCheck(links, client):
    with concurrent.futures.ThreadPoolExecutor(max_workers=5) as executor:
        futures = [executor.submit(asyncThreadRunner, link, headers, client) for link in links]
        return [future.result() for future in as_completed(futures)]
    
    
def get_pages(domain):
    pages = int( httpx.get("https://" + domain + "/wp-json/wp/v2/posts", headers=headers).headers["X-WP-TotalPages"] )
    return pages

def getLinks(rendered_content):
    soup = bs4.BeautifulSoup(rendered_content, "html.parser")
    return [(link["href"], link.text) for link in soup("a") if "href" in link.attrs]

async def fetch_posts(url):
    async with httpx.AsyncClient() as client:
        response = await client.get(url, headers=headers)
        if response.status_code != 200:
            print(f"Error: {response.status_code}")
            return

        data = response.json()
        for post in data:
            post_link = post['link']
            print("Post Link:", post_link, end=" ")

            # Check the status of the post link
            link_response = await client.head(post_link, headers=headers)
            print(link_response.status_code)
            
            # Get links from post content
            links = getLinks(post['content']['rendered'])
            executeBrokenLinkCheck(links, client)

async def main():
    base_url = sys.argv[1]
    
    pages = get_pages(base_url)
    
    print(pages)

    # Iterate through 10 pages, assuming 10 posts per page
    for page in range(1, pages+1):
        url = f"https://{base_url}/wp-json/wp/v2/posts?page={page}"
        print(url)
        print(f"Fetching posts from page {page}...")
        await fetch_posts(url)

asyncio.run(main())
错误信息

运行时出现如下错误:

File "C:\python\lib\asyncio\base_events.py", line 515, in _check_closed
    raise RuntimeError('Event loop is closed')
RuntimeError: Event loop is closed

后来尝试在asyncThreadRunner函数中添加asyncio.set_event_loop_policy(asyncio.WindowsSelectorEventLoopPolicy()),虽然解决了事件循环关闭的错误,但脚本会中途暂停且无法恢复,即使按Ctrl+C也无法退出,只能通过任务管理器手动结束。

问题分析与解决方案

核心问题

你当前的代码混合了线程池与异步IO,导致两个关键问题:

  1. httpx.AsyncClient绑定到创建它的事件循环,不能跨线程使用,强行在子线程调用会引发事件循环冲突。
  2. 在子线程中重复创建事件循环(asyncio.run),导致资源占用异常,最终出现脚本卡死。

修复方案:完全异步化实现

去掉线程池,改用asyncio.gather实现并发链接检测,统一用异步IO处理所有请求,同时完善CSV写入逻辑:

import httpx
import asyncio
import bs4
import sys
import csv
from typing import List, Tuple

headers = {
    'accept': 'text/html,application/xhtml+xml,application/xml;q=0.9,image/avif,image/webp,image/apng,*/*;q=0.8,application/signed-exchange;v=b3;q=0.7',
    'accept-language': 'en-US,en-GB;q=0.9,en;q=0.8',
    'cache-control': 'max-age=0',
    'dnt': '1',
    'priority': 'u=0, i',
    'sec-fetch-dest': 'document',
    'sec-fetch-mode': 'navigate',
    'upgrade-insecure-requests': '1',
    'User-Agent': 'Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/605.1.15 (KHTML, like Gecko) Version/16.0 Safari/605.1.15 Edg/126.0.0.0'
}

# 存储死链的全局列表
broken_links: List[Tuple[str, str, str, str]] = []

async def check_link(link: Tuple[str, str], post_link: str, client: httpx.AsyncClient):
    link_url, link_text = link
    print(f"\t Content Link: {link_url}", end=" ")
    try:
        # 优先用HEAD请求,失败则降级为GET(部分服务器不支持HEAD)
        response = await client.head(link_url, headers=headers, follow_redirects=True)
        status_code = response.status_code
        print(status_code)
        if status_code >= 400:
            broken_links.append((post_link, link_url, link_text, str(status_code)))
    except Exception as errh:
        error_msg = str(errh)
        print(f"{error_msg}")
        broken_links.append((post_link, link_url, link_text, error_msg))

async def execute_broken_link_check(links: List[Tuple[str, str]], post_link: str, client: httpx.AsyncClient):
    # 用信号量控制并发数,避免触发反爬
    semaphore = asyncio.Semaphore(10)
    async def bounded_check(link):
        async with semaphore:
            await check_link(link, post_link, client)
    tasks = [bounded_check(link) for link in links]
    await asyncio.gather(*tasks)

async def get_pages(domain: str, client: httpx.AsyncClient) -> int:
    response = await client.get(f"https://{domain}/wp-json/wp/v2/posts", headers=headers)
    response.raise_for_status()
    return int(response.headers["X-WP-TotalPages"])

def get_links(rendered_content: str) -> List[Tuple[str, str]]:
    soup = bs4.BeautifulSoup(rendered_content, "html.parser")
    return [(link["href"], link.text.strip()) for link in soup.find_all("a", href=True)]

async def fetch_posts(url: str, client: httpx.AsyncClient):
    response = await client.get(url, headers=headers)
    if response.status_code != 200:
        print(f"Error fetching posts: {response.status_code}")
        return

    data = response.json()
    for post in data:
        post_link = post['link']
        print(f"Post Link: {post_link}", end=" ")

        # 检测文章自身链接状态
        try:
            post_response = await client.head(post_link, headers=headers, follow_redirects=True)
            print(post_response.status_code)
            if post_response.status_code >= 400:
                broken_links.append((post_link, post_link, "Post自身链接", str(post_response.status_code)))
        except Exception as err:
            error_msg = str(err)
            print(error_msg)
            broken_links.append((post_link, post_link, "Post自身链接", error_msg))
        
        # 提取文章内的链接并批量检测
        links = get_links(post['content']['rendered'])
        if links:
            await execute_broken_link_check(links, post_link, client)

async def save_to_csv(filename: str = "broken_links.csv"):
    with open(filename, mode='w', newline='', encoding='utf-8') as file:
        writer = csv.writer(file)
        writer.writerow(["来源文章链接", "死链URL", "链接文本", "错误信息/状态码"])
        writer.writerows(broken_links)
    print(f"\n死链已保存到 {filename}")

async def main():
    if len(sys.argv) != 2:
        print("用法: python script.py <域名>")
        sys.exit(1)
    
    base_url = sys.argv[1]
    
    # 全局共用一个AsyncClient,复用连接池提升效率
    async with httpx.AsyncClient(
        follow_redirects=True,
        timeout=httpx.Timeout(10.0),
        headers=headers
    ) as client:
        pages = await get_pages(base_url, client)
        print(f"总共 {pages} 页文章")

        # 并发处理所有页面的文章
        tasks = []
        for page in range(1, pages+1):
            url = f"https://{base_url}/wp-json/wp/v2/posts?page={page}"
            print(f"开始处理第 {page} 页: {url}")
            tasks.append(fetch_posts(url, client))
        
        await asyncio.gather(*tasks)
    
    # 保存死链到CSV
    await save_to_csv()

if __name__ == "__main__":
    # Windows下设置事件循环策略,一次性解决兼容问题
    asyncio.set_event_loop_policy(asyncio.WindowsSelectorEventLoopPolicy())
    asyncio.run(main())

关键改动说明

  1. 移除线程池,全异步实现:用asyncio.gather替代ThreadPoolExecutor,所有IO操作在同一个事件循环中处理,避免跨线程冲突。
  2. 共用AsyncClient:全局复用一个httpx.AsyncClient,利用连接池减少握手开销,同时避免跨线程传递客户端的问题。
  3. 并发数控制:添加Semaphore限制并发请求数,避免触发目标网站的反爬机制或Cloudflare拦截。
  4. 完善错误处理:捕获更多异常类型,将死链信息存入列表后统一写入CSV。
  5. 异步获取总页数:把原同步的get_pages改成异步,避免阻塞事件循环。
  6. Windows兼容:在主入口一次性设置事件循环策略,无需在子线程重复操作。

Cloudflare绕过提示

如果HTTPX仍无法绕过Cloudflare,可以尝试:

  • 升级httpx到最新版本,确保支持TLS 1.3等现代协议。
  • 手动传入浏览器中获取的有效Cloudflare验证Cookie。
  • 配合代理IP分散请求来源。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 21:59:53