如何用Common Crawl与Python快速批量获取数千网页内容?
如何通过Common Crawl和Python快速批量获取数千个网页内容?
我手里有数千个来自不同网站的网页链接,想知道有没有办法通过Common Crawl和Python快速获取这些网页的内容。我目前写了一段代码,但处理速度很慢,代码如下:
async def search_cc_index(url): encoded_url = quote_plus(url) index_url = f'{SERVER}{INDEX_NAME}-index?url={encoded_url}&output=json' async with aiohttp.ClientSession() as session: async with session.get(index_url) as response: if response.status == 200: records = (await response.text()).strip().split('\n') return [json.loads(record) for record in records] else: return None async def fetch_page_from_cc(records): async with aiohttp.ClientSession() as session: for record in records: offset, length = int(record['offset']), int(record['length']) s3_url = f'https://data.commoncrawl.org/{record["filename"]}' byte_range = f'bytes={offset}-{offset + length - 1}' async with session.get(s3_url, headers={'Range': byte_range}) as response: if response.status == 206: stream = ArchiveIterator(response.content) for warc_record in stream: if warc_record.rec_type == 'response': return await warc_record.content_stream().read() else: return None return None async def fetch_individual_url(target_url): records = await search_cc_index(target_url) if records: print(f"Found {len(records)} records for {target_url}") content = await fetch_page_from_cc(records) if content: print(f"Successfully fetched content for {target_url}") else: print(f"No records found for {target_url}")
优化方案
1. 复用ClientSession
代码里每个请求都新建ClientSession,会产生额外的连接建立开销。应该全局复用一个Session,在批量任务开始时初始化,结束时关闭。
2. 并发批量处理URL
当前是逐个处理URL,完全没发挥异步的优势。用asyncio.Semaphore控制并发数,配合asyncio.gather同时发起多个请求,既能提速又不会触发Common Crawl的限流。
3. 优化索引查询逻辑
- 过滤出状态码为200的有效记录,避免无效请求;按抓取时间排序,优先取最新的记录。
- 给索引查询和WARC内容请求加上重试机制,应对偶尔的网络波动或服务响应失败。
4. 减少重复操作
- 缓存已查询过的URL结果,避免重复查询相同URL。
- 如果多个URL对应同一个WARC文件,可以一次性下载对应字节范围,批量提取内容(适合URL数量极大的场景)。
优化后的示例代码
import asyncio import json from urllib.parse import quote_plus import aiohttp from warcio.archiveiterator import ArchiveIterator from tenacity import retry, stop_after_attempt, wait_exponential # 配置常量 SERVER = "https://index.commoncrawl.org/" INDEX_NAME = "CC-MAIN-2024-22" # 全局复用ClientSession session = None async def init_session(): global session session = aiohttp.ClientSession() async def close_session(): await session.close() @retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=2, max=10)) async def search_cc_index(url): encoded_url = quote_plus(url) index_url = f'{SERVER}{INDEX_NAME}-index?url={encoded_url}&output=json' async with session.get(index_url) as response: if response.status == 200: records_text = await response.text() records = [json.loads(line) for line in records_text.strip().split('\n') if line] # 过滤200状态的记录,按时间取最新的 valid_records = [r for r in records if r.get('status') == '200'] if valid_records: valid_records.sort(key=lambda x: x.get('timestamp', ''), reverse=True) return valid_records return None @retry(stop=stop_after_attempt(3), wait=wait_exponential(multiplier=1, min=2, max=10)) async def fetch_page_from_cc(record): offset, length = int(record['offset']), int(record['length']) s3_url = f'https://data.commoncrawl.org/{record["filename"]}' byte_range = f'bytes={offset}-{offset + length - 1}' async with session.get(s3_url, headers={'Range': byte_range}) as response: if response.status == 206: stream = ArchiveIterator(response.content) for warc_record in stream: if warc_record.rec_type == 'response': return await warc_record.content_stream().read() return None async def fetch_individual_url(target_url): records = await search_cc_index(target_url) if records: print(f"找到 {len(records)} 条有效记录 for {target_url}") # 只取最新的一条记录抓取 content = await fetch_page_from_cc(records[0]) if content: print(f"成功获取 {target_url} 的内容") # 这里可以添加保存内容到文件/数据库的逻辑 return content else: print(f"无法获取 {target_url} 的内容") else: print(f"{target_url} 没有找到匹配的记录") return None async def batch_fetch_urls(url_list, max_concurrent=20): # 控制并发数,避免被限流 semaphore = asyncio.Semaphore(max_concurrent) async def bounded_fetch(url): async with semaphore: return await fetch_individual_url(url) await init_session() try: tasks = [bounded_fetch(url) for url in url_list] await asyncio.gather(*tasks) finally: await close_session() # 使用示例 if __name__ == "__main__": url_list = ["https://example.com", "https://test.com"] # 替换成你的数千个URL asyncio.run(batch_fetch_urls(url_list))
内容的提问来源于stack exchange,提问作者ALTAF HUSSAIN
相关产品推荐
相关产品推荐

