Python asyncio实战:用10行代码书写高并发视频信息采集器

2026-08-31 0 204

最近朋友开了个短视频数据分析的小工作室,让我帮忙弄个工具:从几个公开的视频网站抓取视频标题、播放数和发布时间,用来做行业报告。视频网站接口数量不多,但单个接口返回慢得很,最少也要几百毫秒。一开始我用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的主要优势不是代码多炫酷,而是把并发模型简化了,让你不用手动创建线程或者用回调,只要顺序书写异步代码就行了。这篇文章就说这么多,代码能直接用,有问题评论区见。

Python asyncio实战:用10行代码书写高并发视频信息采集器
收藏 (0) 打赏

感谢您的支持,我会继续努力的!

打开微信/支付宝扫一扫,即可进行扫码打赏哦,分享从这里开始,精彩与您同在
点赞 (0)

版权声明:
本站资源有的来自互联网收集整理,本站纯免费分享提供学习使用,如果侵犯了您的合法权益,请发送邮件1506151422@qq.com联系,将会及时下架删除。
本站资源仅供研究、学习交流之用,免费开源项目不代表完全可商用,若商业用途请先咨询开发企业能否商用,否则产生的一切后果将由下载用户自行承担。
原创板块未经允许不得转载,否则将追究法律责任。

淘吗网 python Python asyncio实战:用10行代码书写高并发视频信息采集器 https://www.taomawang.com/server/python/2678.html

下一篇:

已经没有下一篇了!

常见问题

相关文章

猜你喜欢
发表评论
暂无评论
官方客服团队

为您解决烦忧 - 24小时在线 专业服务