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

生产脚本获取千余设备:httpx/ThreadPoolExecutor/aiohttp选型与最佳实践

批量抓取1214台设备数据:httpx异步、ThreadPoolExecutor+requests还是aiohttp?

方案对比与选择建议

1. ThreadPoolExecutor + requests

  • 适用场景:团队不熟悉异步编程,或任务量不大、对性能要求不极致的场景。
  • 优势:同步代码逻辑简单,调试方便,requests生态成熟,处理常规并发(16-64线程)时稳定性强。
  • 劣势:线程本身有内存开销,大量线程切换会消耗CPU资源,1214个任务的总耗时会比异步方案长。

2. httpx 异步

  • 适用场景:追求性能且希望降低学习成本的IO密集型任务(比如批量HTTP请求)。
  • 优势:API设计和requests高度一致,迁移成本极低;协程开销远小于线程,能支持更高并发量;同时支持同步/异步两种模式,灵活度高。
  • 劣势:需要掌握asyncio基本概念(如await、事件循环),但上手门槛远低于aiohttp。

3. aiohttp

  • 适用场景:超大规模异步任务(上万级请求),或需要精细控制连接池、超时、重试等底层细节的场景。
  • 优势:老牌异步HTTP库,底层优化完善,性能稳定,社区资源丰富。
  • 劣势:API和requests差异较大,学习成本较高,需要手动管理会话和并发控制的更多细节。

1214台设备抓取的最佳实践

针对1214个设备的批量请求,优先推荐httpx异步或aiohttp,协程的低开销能显著缩短总耗时。以下是优化后的代码示例:

优化后的httpx异步方案(带并发控制)

import asyncio
import httpx

devices = ["device_1", ...]  # 1214个设备

async def fetch_single_device(client, device, semaphore):
    # 用信号量限制并发数,避免触发服务器限流
    async with semaphore:
        try:
            # 设置超时时间,防止请求卡住
            response = await client.get(
                f"https://myserver.com/api/devices/{device}",
                timeout=httpx.Timeout(10.0)
            )
            response.raise_for_status()  # 主动抛出HTTP错误
            return response.json()
        except httpx.HTTPError as e:
            print(f"设备 {device} 抓取失败: {str(e)}")
            return None

async def fetch_all_devices():
    # 根据服务器承受能力调整并发数,建议50-100之间
    semaphore = asyncio.Semaphore(60)
    async with httpx.AsyncClient() as client:
        # 生成所有任务
        tasks = [fetch_single_device(client, dev, semaphore) for dev in devices]
        # 批量执行并收集结果
        results = await asyncio.gather(*tasks)
        # 过滤失败的请求
        successful_results = [res for res in results if res is not None]
        print(f"成功抓取 {len(successful_results)}/{len(devices)} 台设备数据")

if __name__ == "__main__":
    asyncio.run(fetch_all_devices())

优化后的ThreadPoolExecutor方案(带结果处理)

import requests
from concurrent.futures import ThreadPoolExecutor, as_completed

devices = ["device_1", ...]  # 1214个设备

def fetch_single_device(session, device):
    try:
        response = session.get(
            f"https://myserver.com/api/devices/{device}",
            timeout=10
        )
        response.raise_for_status()
        return response.json()
    except requests.exceptions.RequestException as e:
        print(f"设备 {device} 抓取失败: {str(e)}")
        return None

def fetch_all_devices():
    successful_results = []
    # 复用Session,减少TCP连接建立开销
    with requests.Session() as session:
        # 根据本地资源调整线程数,建议32-64
        with ThreadPoolExecutor(max_workers=40) as executor:
            futures = {executor.submit(fetch_single_device, session, dev): dev for dev in devices}
            # 按完成顺序处理结果
            for future in as_completed(futures):
                result = future.result()
                if result is not None:
                    successful_results.append(result)
        print(f"成功抓取 {len(successful_results)}/{len(devices)} 台设备数据")

if __name__ == "__main__":
    fetch_all_devices()

通用注意事项

  • 并发控制:无论哪种方案,都不要一次性发起所有请求,必须通过信号量(异步)或线程数(线程池)限制并发量,避免触发服务器限流或导致服务不可用。
  • 错误处理:必须捕获HTTP异常、超时等错误,避免单个请求失败拖垮整个任务。
  • 会话复用:复用requests.Session、httpx.AsyncClient或aiohttp.ClientSession,减少TCP连接建立的开销,提升请求效率。
  • 超时设置:给每个请求设置合理的超时时间,防止因单个请求卡住导致任务整体耗时过长。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 19:07:04