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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 13:35:11