分布式任务调度 —— 用 Celery + Playwright 构建高并发爬虫集群
·
分布式任务调度 —— 用 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() 包裹异步逻辑;
- 每个任务独立启动浏览器,避免上下文污染;
- 自动重试机制应对临时网络波动。
✅ 步骤四:启动服务
- 启动 Redis
- 启动 Celery
# 终端 1:启动 worker(可开多个)
celery -A celery_app worker --loglevel=info --concurrency=4
# --concurrency=4 表示每个 worker 进程同时处理 4 个任务
# 实际并发数 = worker 数 × concurrency
- 生产任务
# 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 查看实时任务状态、成功率、耗时等。
⚠️ 注意事项 & 最佳实践
- 资源控制:别让浏览器吃光内存
- 每个 Playwright 实例约占用 100~300MB 内存;
- 建议 --concurrency=2~4(根据机器配置调整);
- 使用 headless=True 减少资源消耗。
- 登录态共享 vs 独立
- 如果所有任务用同一个账号,可共享 bilibili_login_state.json;
- 如果需多账号轮换,可将 state 文件路径作为任务参数传入。
- 结果持久化
- 任务返回的数据默认存在 Redis,长期运行会撑爆内存。建议:
- 将结果写入 MySQL / MongoDB;
- 或使用 @app.task(…, ignore_result=True) 忽略结果。
- 优雅关闭
- 按 Ctrl+C 时,Celery 会等待当前任务完成再退出,确保数据不丢失。
深入架构设计 —— 分布式任务调度系统
📝 架构设计详解
- 生产者-消费者模型
我们的架构基于生产者-消费者模式,其中:
- 生产者负责生成任务并将其推送到消息队列(Redis)。
- 消费者(即Celery Workers)从消息队列中拉取任务,并执行相应的处理逻辑(Playwright抓取网页内容)。
这种设计的一个关键优势是解耦了任务的生成与处理,允许我们在不影响整体系统的情况下独立优化这两部分。
- 水平扩展能力
通过增加更多的Celery Worker实例,我们可以轻松地提升系统的处理能力。每个Worker都是相对独立的,可以运行在不同的机器上,这使得系统具有良好的水平扩展性。
实现方式:
- 增加Worker数量:可以通过在同一台或多台服务器上启动多个Worker进程来实现。
- 负载均衡:由消息队列(如Redis)自动分配任务给空闲的Worker,无需人工干预。
- 容错机制
Celery提供了强大的错误重试和结果持久化功能,这对于网络请求不稳定或需要保证数据完整性的场景尤为重要。
- 自动重试:当某个任务失败时(例如网络超时),Celery可以自动重新尝试执行该任务。
- 结果存储:通过配置backend(如Redis),Celery可以保存每个任务的执行结果,便于后续查询或审计。
架构流程图
- Task Producer:表示任务的生产者,负责生成需要执行的任务。
- Redis Broker:作为消息队列的角色,用来存储和分发任务。
- Workers Cluster:代表由多个Celery Worker组成的集群,每个Worker都可以独立处理任务。
- Playwright Instances:每个Worker都会启动一个Playwright实例来执行具体的网页抓取任务。
- Target Web Servers:表示Playwright抓取的目标Web服务器。
更多推荐
所有评论(0)