最近朋友开了个短视频数据分析的小工作室,让我帮忙弄个工具:从几个公开的视频网站抓取视频标题、播放数和发布时间,用来做行业报告。视频网站接口数量不多,但单个接口返回慢得很,最少也要几百毫秒。一开始我用requests一条条请求,抓100个视频要2分多钟,而且网络一抖就卡死。后来我用asyncio重写了采集逻辑,同样的任务秒级就完成了。这篇文章把我踩过的坑和最终能用的代码写出来。
先搞定asyncio的基本理解
asyncio是Python的异步IO库,核心就是事件循环。你可以把任务想象成一堆烧水壶,事件循环就是那个同时看着所有水壶的人,哪个水壶响了就去处理哪个,而不是等一个烧开再去烧下一个。这特别适合网络IO密集型的任务,因为等网络响应的时候CPU完全没事干。
咱们用的关键东西有两个:aiohttp发起异步HTTP请求,asyncio管理协程并发。
先装好依赖:
pip install aiohttp
不需要装别的,Python版本3.8以上就行,我用的3.12。
最傻的异步版本:一个一个来
先别急着并发,我们写一个最基础的异步采集单条视频信息的函数。假设有个接口返回JSON:
import asyncio
import aiohttp
async def fetch_video_info(session, video_id):
url = f"https://api.example.com/video/{video_id}"
async with session.get(url, timeout=5) as resp:
data = await resp.json()
return {
"id": video_id,
"title": data["title"],
"play_count": data["play_count"],
"publish_time": data["publish_time"],
}
这里用了async with自动管理会话,await resp.json()把JSON解析也挂起,等数据返回。注意,用requests会阻塞事件循环,千万别混用。
如果要采集单个视频,这样就能跑:
async def main():
async with aiohttp.ClientSession() as session:
info = await fetch_video_info(session, "1001")
print(info)
asyncio.run(main())
但真要这么干,跟requests没啥区别。我们需要并发。
用gather实现真正的并发
并发的最简单写法是用asyncio.gather,把一堆协程扔进去等全部完成。比如要抓5个视频:
async def main():
async with aiohttp.ClientSession() as session:
video_ids = ["1001", "1002", "1003", "1004", "1005"]
tasks = [fetch_video_info(session, vid) for vid in video_ids]
results = await asyncio.gather(*tasks, return_exceptions=True)
for result in results:
if isinstance(result, Exception):
print("有任务失败了:", result)
else:
print(result)
注意return_exceptions=True,这样即使某个请求挂了,也不会让整个程序崩掉,而是返回一个异常对象。这是必须的,因为网络请求总有失败的时候。
不过上面这个写法如果视频数量很大,比如1000个,同时发起1000个并发请求,很可能会被目标网站封IP。所以需要限流。
限流:用信号量控制并发数
给自己的采集器加个并发上限,防止把对方服务器搞崩,也防止自己IP被ban。asyncio里有个Semaphore,可以限制同时运行的协程数量。
async def fetch_video_info(session, video_id, semaphore):
async with semaphore:
# 这里就是需要限制并发的部分
url = f"https://api.example.com/video/{video_id}"
async with session.get(url, timeout=5) as resp:
data = await resp.json()
return {...}
然后在main里创建信号量,比如同时只跑20个:
async def main():
semaphore = asyncio.Semaphore(20)
async with aiohttp.ClientSession() as session:
tasks = [fetch_video_info(session, vid, semaphore) for vid in video_ids]
results = await asyncio.gather(*tasks, return_exceptions=True)
这样就不用担心并发数太高了。
处理重试和超时
实际采集的时候,某个接口偶尔会超时,这时候最好能重试几次。用asyncio重试不能直接用time.sleep,要用await asyncio.sleep()。我在采集函数里面加了最简单的重试逻辑:
async def fetch_video_info(session, video_id, semaphore):
async with semaphore:
url = f"https://api.example.com/video/{video_id}"
for attempt in range(3):
try:
async with session.get(url, timeout=5) as resp:
if resp.status != 200:
raise aiohttp.ClientError(f"HTTP {resp.status}")
data = await resp.json()
return {...}
except (aiohttp.ClientError, asyncio.TimeoutError) as e:
if attempt == 2:
raise # 重试3次还失败,抛给上层
wait_time = 0.5 * (2 ** attempt)
print(f"视频{video_id}第{attempt+1}次失败,{wait_time}秒后重试")
await asyncio.sleep(wait_time)
这里为了不把文章搞太复杂,重试间隔用了指数退避,但没加抖动。真要用于生产环境,建议加一点随机抖动,防止多个请求同时重试造成波动。
一个完整的可运行示例
把上面的东西拼起来,写成一个模块。这个模块可以从文件里读取视频ID,简单统计用时和结果,保存到JSON文件。
import asyncio
import aiohttp
import json
import time
async def fetch_video_info(session, video_id, semaphore):
async with semaphore:
url = f"https://api.example.com/video/{video_id}"
for attempt in range(3):
try:
async with session.get(url, timeout=5) as resp:
if resp.status != 200:
raise aiohttp.ClientError(f"HTTP {resp.status}")
data = await resp.json()
return {
"id": video_id,
"title": data["title"],
"play_count": data["play_count"],
"publish_time": data["publish_time"],
}
except (aiohttp.ClientError, asyncio.TimeoutError) as e:
if attempt == 2:
raise
wait_time = 0.5 * (2 ** attempt)
await asyncio.sleep(wait_time)
async def main(video_ids):
semaphore = asyncio.Semaphore(20)
async with aiohttp.ClientSession() as session:
tasks = [fetch_video_info(session, vid, semaphore) for vid in video_ids]
results = await asyncio.gather(*tasks, return_exceptions=True)
success = []
failed = []
for i, result in enumerate(results):
if isinstance(result, Exception):
failed.append({"id": video_ids[i], "error": str(result)})
else:
success.append(result)
return success, failed
if __name__ == "__main__":
ids = [f"{i:04d}" for i in range(1, 101)] # 模拟100个视频ID
start = time.perf_counter()
ok, bad = asyncio.run(main(ids))
elapsed = time.perf_counter() - start
print(f"成功 {len(ok)} 条,失败 {len(bad)} 条,耗时 {elapsed:.2f} 秒")
with open("result.json", "w", encoding="utf-8") as f:
json.dump(ok, f, ensure_ascii=False, indent=2)
if bad:
with open("error.json", "w", encoding="utf-8") as f:
json.dump(bad, f, ensure_ascii=False, indent=2)
这代码在本地跑的时候,我故意把视频接口地址换成一个慢接口,每个请求要800ms。用100个视频测试,requests同步版耗时80秒左右,asyncio并发20个,耗时5秒左右,提速16倍。如果再加大并发到50,能达到3秒多,但对方服务器可能会拉黑我,所以还是控制在20比较稳妥。
实际采集中的几个坑
1. aiohttp的session要复用
不要在每个协程里都创建新的ClientSession。Session内部维护了连接池,全局共享一个就好,否则会创建大量socket,性能反而更差。
2. 注意resp.release()
如果你用resp.read()或者resp.json(),连接会被自动释放。但如果你拿不到响应体,比如超时了,连接可能没有被正确归还连接池。用async with session.get(...)可以保证连接最终会释放。
3. 用绝对路径写文件
如果你把那个脚本丢到crontab运行,相对路径可能会找不到文件。我都是写死绝对路径,这个不算asyncio的问题,但容易忽略。
4. 爬太快会被封IP
如果你从本机直接跑1000个并发,目标网站的防火墙很容易识别出来。我用了两个方法:设置信号量限制并发;设置每个请求的随机小延迟,比如await asyncio.sleep(random.uniform(0.1, 0.5))。这样做虽然慢一点,但安全。
能不能更优雅一点?用gather还是TaskGroup
Python 3.11以后提供了asyncio.TaskGroup,它的写法更结构化,而且遇到异常会自动取消还没完成的任务。如果你用的是Python 3.11+,可以这样写:
async def main(video_ids):
semaphore = asyncio.Semaphore(20)
results = {}
async with aiohttp.ClientSession() as session:
async with asyncio.TaskGroup() as tg:
tasks = [tg.create_task(fetch_with_sem(session, vid, semaphore)) for vid in video_ids]
# 这时所有任务已完成或抛异常
for task in tasks:
results[task.get_name()] = task.result()
但TaskGroup的异常处理跟gather不太一样,如果某个任务抛了异常,TaskGroup会直接抛出来,并且自动取消其他任务。这有时候并不是我们想要的,毕竟一个视频失败不应该影响其他视频采集。所以在采集场景下,我还是倾向于用gather(return_exceptions=True)。
结尾:以后这种事可以怎么干
如果你只是临时采个几百条数据,随便写写就得了。但要是长期跑,我建议把视频ID存到Redis或者数据库里,然后分段拉取,每批100个,采集完成后记录状态。关于asyncio的坑,你踩过几个之后就能形成自己的避坑指南了。
反正我用这套代码已经稳定跑了一周,每天抓两万条数据,没出过幺蛾子。我觉得asyncio的主要优势不是代码多炫酷,而是把并发模型简化了,让你不用手动创建线程或者用回调,只要顺序书写异步代码就行了。这篇文章就说这么多,代码能直接用,有问题评论区见。

