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,导致两个关键问题:
httpx.AsyncClient绑定到创建它的事件循环,不能跨线程使用,强行在子线程调用会引发事件循环冲突。- 在子线程中重复创建事件循环(
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())
关键改动说明
- 移除线程池,全异步实现:用
asyncio.gather替代ThreadPoolExecutor,所有IO操作在同一个事件循环中处理,避免跨线程冲突。 - 共用AsyncClient:全局复用一个
httpx.AsyncClient,利用连接池减少握手开销,同时避免跨线程传递客户端的问题。 - 并发数控制:添加
Semaphore限制并发请求数,避免触发目标网站的反爬机制或Cloudflare拦截。 - 完善错误处理:捕获更多异常类型,将死链信息存入列表后统一写入CSV。
- 异步获取总页数:把原同步的
get_pages改成异步,避免阻塞事件循环。 - Windows兼容:在主入口一次性设置事件循环策略,无需在子线程重复操作。
Cloudflare绕过提示
如果HTTPX仍无法绕过Cloudflare,可以尝试:
- 升级httpx到最新版本,确保支持TLS 1.3等现代协议。
- 手动传入浏览器中获取的有效Cloudflare验证Cookie。
- 配合代理IP分散请求来源。
内容的提问来源于stack exchange,提问作者Salman Khan
相关产品推荐
相关产品推荐

