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

如何将基于requests的船舶数据爬取代码改造为asyncio+aiohttp异步版本提速

同步Requests船舶爬取代码转异步(asyncio+aiohttp)

原代码基于requests同步请求MarineTraffic船舶数据,因串行等待响应导致爬取速度极慢。以下是改造后的异步版本,利用asyncio和aiohttp实现并发请求,大幅提升爬取效率。

改造后的完整代码

import asyncio
import aiohttp
from time import perf_counter

# 文件写入锁,避免异步并发写入冲突
write_lock = asyncio.Lock()

async def get_ship_position(session, ship_id):
    url = "https://www.marinetraffic.com/en/vesselDetails/vesselInfo/shipid:{}".format(ship_id)
    
    headers = {
        "accept": "application/json",
        "accept-encoding": "gzip, deflate",
        "user-agent": "Mozilla/5.0",
        "x-requested-with": "XMLHttpRequest"
    }

    async with session.get(url, headers=headers) as response:
        response.raise_for_status()
        return await response.json()

async def process_ship(session, ship_id):
    try:
        data = await get_ship_position(session, ship_id)
        # 加锁写入成功数据
        async with write_lock:
            with open("marinetraffic.txt", "a", encoding="utf-8") as bos:
                line = "{}\t{}\t{}\t{}\t{}\t{}\t{}\t{}x{}\t{}\t{}\t{}\t{}\t{}\t{}".format(
                    data["mmsi"], data["imo"], data["name"], data["nameAis"],
                    data["type"], data["typeSpecific"], data["yearBuilt"],
                    data["length"], data["breadth"], data["callsign"],
                    data["country"], data["deadweight"], data["grossTonnage"],
                    data["homePort"], data["status"]
                )
                print(line, file=bos)
        print(ship_id, "Yazdı")
    except Exception as e:
        print(ship_id, "Hata")
        # 加锁写入错误数据
        async with write_lock:
            with open("marinetraffichata.txt", "a", encoding="utf-8") as hata:
                print("Hata", ship_id, file=hata)

async def main():
    start = perf_counter()
    # 创建aiohttp客户端会话,复用连接池
    async with aiohttp.ClientSession() as session:
        # 生成所有爬取任务
        tasks = [process_ship(session, ship_id) for ship_id in range(7551, 10000)]
        # 并发执行所有任务
        await asyncio.gather(*tasks)
    stop = perf_counter()
    print("çalışılan süre:", stop - start, "saniye")

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

关键改动说明

  • 异步请求替换:用aiohttp.ClientSession替代requests,通过await session.get()发起异步请求,避免串行等待响应的时间浪费。
  • 异步函数改造:所有涉及IO操作的函数改为async def,用await标记需要等待的异步操作(请求发送、响应解析)。
  • 并发任务管理:生成全部爬取任务的列表,通过asyncio.gather()一次性并发执行,同时发起多个HTTP请求。
  • 文件写入安全:添加asyncio.Lock(),确保多个异步任务同时写入文件时不会出现数据错乱,解决并发写入冲突问题。
  • 连接池复用:aiohttp.ClientSession自动复用TCP连接,减少重复建立连接的开销,进一步提升效率。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 18:20:43