分布式任务调度 —— 用 Celery + Playwright 构建高并发爬虫集群-- pd的爬虫笔记

环境搭建

✅ 步骤一:安装依赖

pip install celery redis 

✅ 步骤二:配置 Celery(celery_app.py)

# celery_app.py
from celery import Celery

# 使用 Redis 作为 broker 和 backend
app = Celery(
    "bilibili_crawler",
    broker="redis://localhost:6379/0",
    backend="redis://localhost:6379/0",
)

# 全局配置
app.conf.update(
    task_serializer="json",
    accept_content=["json"],
    result_serializer="json",
    timezone="Asia/Shanghai",
    enable_utc=False,
    task_routes={
        "tasks.fetch_user_profile": {"queue": "bilibili_high_priority"},
    },
)

✅ 步骤三:编写 Playwright 任务(tasks.py)

这个只是示例,实际任务请根据业务需求编写。

# tasks.py
import asyncio
import os
from playwright.async_api import async_playwright
from celery_app import app
from stealth_config import create_stealth_browser  # 复用你的 stealth 配置

@app.task(bind=True, autoretry_for=(Exception,), retry_kwargs={"max_retries": 3, "countdown": 5})
def fetch_user_profile(self, uid: int):
    """
    Celery 任务:抓取 B站用户主页
    注意:Celery 默认不支持 async,需手动运行 asyncio
    """
    return asyncio.run(_fetch_user_profile_async(uid))

async def _fetch_user_profile_async(uid: int):
    state_file = "bilibili_login_state.json"
    if not os.path.exists(state_file):
        raise FileNotFoundError("登录态文件缺失,请先运行 save_login_state.py")

    # 启动带登录态的 stealth 浏览器
    playwright, browser, context = await create_stealth_browser(
        headless=True,
        storage_state=state_file
    )
    
    try:
        page = await context.new_page()
        url = f"https://space.bilibili.com/{uid}"
        await page.goto(url, timeout=30000)

        # 验证是否登录(避免被重定向到登录页)
        if "passport.bilibili.com" in page.url:
            raise RuntimeError("登录态失效")

        # 等待关键元素加载(如用户名)
        await page.wait_for_selector(".username", timeout=10000)
        username = await page.text_content(".username")

        # 提取更多数据...
        follower_count = await page.text_content(".n-fans") or "0"

        return {
            "uid": uid,
            "username": username.strip(),
            "followers": follower_count.replace("粉丝", "").strip(),
            "url": url,
            "status": "success"
        }

    except Exception as e:
        return {
            "uid": uid,
            "error": str(e),
            "status": "failed"
        }
    finally:
        await context.close()
        await browser.close()
        await playwright.stop()

⚠️ 关键点:

  • Celery 任务函数必须是 同步函数,所以用 asyncio.run() 包裹异步逻辑;
  • 每个任务独立启动浏览器,避免上下文污染;
  • 自动重试机制应对临时网络波动。

✅ 步骤四:启动服务

  1. 启动 Redis
  2. 启动 Celery
# 终端 1:启动 worker(可开多个)
celery -A celery_app worker --loglevel=info --concurrency=4

# --concurrency=4 表示每个 worker 进程同时处理 4 个任务
# 实际并发数 = worker 数 × concurrency
  1. 生产任务
# submit_tasks.py
from tasks import fetch_user_profile

# 批量提交 1000 个用户 ID
uids = list(range(1, 1001))
for uid in uids:
    task = fetch_user_profile.delay(uid)  # 异步提交
    print(f"📤 提交任务: UID={uid}, TaskID={task.id}")

📊 查看结果 & 监控

获取单个任务结果

from celery.result import AsyncResult
from celery_app import app

result = AsyncResult("your-task-id", app=app)
print(result.get(timeout=10))  # 阻塞等待结果

使用 Flower 监控面板(可选)

pip install flower
celery -A celery_app flower --port=5555

访问 http://localhost:5555 查看实时任务状态、成功率、耗时等。

⚠️ 注意事项 & 最佳实践

  1. 资源控制:别让浏览器吃光内存
    • 每个 Playwright 实例约占用 100~300MB 内存;
    • 建议 --concurrency=2~4(根据机器配置调整);
    • 使用 headless=True 减少资源消耗。
  2. 登录态共享 vs 独立
    • 如果所有任务用同一个账号,可共享 bilibili_login_state.json;
    • 如果需多账号轮换,可将 state 文件路径作为任务参数传入。
  3. 结果持久化
    • 任务返回的数据默认存在 Redis,长期运行会撑爆内存。建议:
    • 将结果写入 MySQL / MongoDB;
    • 或使用 @app.task(…, ignore_result=True) 忽略结果。
  4. 优雅关闭
    • 按 Ctrl+C 时,Celery 会等待当前任务完成再退出,确保数据不丢失。

深入架构设计 —— 分布式任务调度系统

📝 架构设计详解

  1. 生产者-消费者模型

我们的架构基于生产者-消费者模式,其中:

  • 生产者负责生成任务并将其推送到消息队列(Redis)。
  • 消费者(即Celery Workers)从消息队列中拉取任务,并执行相应的处理逻辑(Playwright抓取网页内容)。
    这种设计的一个关键优势是解耦了任务的生成与处理,允许我们在不影响整体系统的情况下独立优化这两部分。
  1. 水平扩展能力

通过增加更多的Celery Worker实例,我们可以轻松地提升系统的处理能力。每个Worker都是相对独立的,可以运行在不同的机器上,这使得系统具有良好的水平扩展性。

实现方式:

  • 增加Worker数量:可以通过在同一台或多台服务器上启动多个Worker进程来实现。
  • 负载均衡:由消息队列(如Redis)自动分配任务给空闲的Worker,无需人工干预。
  1. 容错机制

Celery提供了强大的错误重试和结果持久化功能,这对于网络请求不稳定或需要保证数据完整性的场景尤为重要。

  • 自动重试:当某个任务失败时(例如网络超时),Celery可以自动重新尝试执行该任务。
  • 结果存储:通过配置backend(如Redis),Celery可以保存每个任务的执行结果,便于后续查询或审计。

架构流程图

Target Web Servers

Playwright Instances

Workers Cluster

Message Queue

Push Task to Redis

Task Producer

Redis Broker

Worker 1

Worker 2

Worker N

Playwright Instance 1

Playwright Instance 2

Playwright Instance N

Web Server 1

Web Server 2

Web Server N

  • Task Producer:表示任务的生产者,负责生成需要执行的任务。
  • Redis Broker:作为消息队列的角色,用来存储和分发任务。
  • Workers Cluster:代表由多个Celery Worker组成的集群,每个Worker都可以独立处理任务。
  • Playwright Instances:每个Worker都会启动一个Playwright实例来执行具体的网页抓取任务。
  • Target Web Servers:表示Playwright抓取的目标Web服务器。
Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐