上一版只把客户端写出来、验证了它单独可用,**但没有接进任何地方** —— 所以你建的任务跑起来 仍然在调爬虫子进程,报的仍然是那句自己编的「account blocked」。这一步把它接上。 * runner 的 Phase 2 按平台分岔:dy 走进程内 HTTP 客户端(douyin_fetch),其余平台照旧走 爬虫子进程。抖音那条不再起 Playwright、不再构造那串自相矛盾的浏览器指纹参数。 * 新增 douyin_fetch:把采到的东西写成 store/douyin 那套 jsonl 形状 —— **ingest 完全不知道 数据是从哪来的**,重采样/差分/事件/通知/报表全都照旧,一个字没改。 * 失败不再假装:一条都没采到就以非零退出码 + **真实原因**交给 ingest,落成 「采集进程异常退出(code=1):…」。绝不会再掉进「疑似登录失效」那个分支。 * 已知作品列表接口(aweme/post)被抖音单独加了真校验(200 + 空 body),所以加了退化: 拿不到列表就用库里已知的 aweme_id 逐条走 detail 刷新。**边界是:已知作品的指标能继续 更新,新作品发现不了** —— 这个边界会以一条 warning 日志留下痕迹,不让它看起来一切正常。 * 顺带给客户端补上 video_detail(实测可用:200 / 45425 字节),退化路径靠它。 测试 +6:产物目录与文件名、评论文件即使为空也要建(ingest 靠它区分「没评论」和 「什么都没抓到」)、重复作品只写一次、列表被挡时的退化、彻底失败仍写出产物与原因、 评论失败不连累作品。
356 lines
15 KiB
Python
356 lines
15 KiB
Python
# -*- coding: utf-8 -*-
|
||
# Copyright (c) 2025 [email protected]
|
||
#
|
||
# This file is part of MediaCrawler project.
|
||
# Repository: https://github.com/NanmiCoder/MediaCrawler/blob/main/api/monitor/runner.py
|
||
# GitHub: https://github.com/NanmiCoder
|
||
# Licensed under NON-COMMERCIAL LEARNING LICENSE 1.1
|
||
#
|
||
# 声明:本代码仅供学习和研究目的使用。使用者应遵守以下原则:
|
||
# 1. 不得用于任何商业用途。
|
||
# 2. 使用时应遵守目标平台的使用条款和robots.txt规则。
|
||
# 3. 不得进行大规模爬取或对平台造成运营干扰。
|
||
# 4. 应合理控制请求频率,避免给目标平台带来不必要的负担。
|
||
# 5. 不得用于任何非法或不当的用途。
|
||
#
|
||
# 详细许可条款请参阅项目根目录下的LICENSE文件。
|
||
# 使用本代码即表示您同意遵守上述原则和LICENSE中的所有条款。
|
||
|
||
"""Execute a single monitoring run: build the command, wait, then ingest.
|
||
|
||
Runs reuse ``CrawlerManager`` so that monitor crawls share the existing
|
||
single-subprocess guarantee and their logs stream to the existing Terminal
|
||
component over the existing log WebSocket.
|
||
"""
|
||
|
||
import asyncio
|
||
import os
|
||
from pathlib import Path
|
||
from typing import Iterable, List, Optional
|
||
|
||
from sqlalchemy import select
|
||
|
||
from tools.time_util import get_current_timestamp
|
||
|
||
from ..schemas import (
|
||
CrawlerStartRequest,
|
||
CrawlerTypeEnum,
|
||
LoginTypeEnum,
|
||
PlatformEnum,
|
||
SaveDataOptionEnum,
|
||
)
|
||
from ..services import crawler_manager
|
||
from . import adapters, app_settings, covers, douyin_fetch, notify
|
||
from .db import get_session
|
||
from .ingest import IngestResult, diagnose_failure, ingest_run
|
||
from .models import (
|
||
MODE_CREATOR,
|
||
RUN_FAILED,
|
||
RUN_PENDING,
|
||
RUN_RUNNING,
|
||
RUN_TIMEOUT,
|
||
MonitorNote,
|
||
MonitorRun,
|
||
MonitorTarget,
|
||
MonitorTask,
|
||
)
|
||
from .settings import get_cookie, mark_cookie_ok
|
||
|
||
PROJECT_ROOT = Path(__file__).parent.parent.parent
|
||
MONITOR_RUNS_DIR = PROJECT_ROOT / "data" / "monitor_runs"
|
||
|
||
# Monitor platform ids align with PlatformEnum's values, but mapping explicitly
|
||
# beats relying on that coincidence.
|
||
_PLATFORM_ENUM = {
|
||
"xhs": PlatformEnum.XHS,
|
||
"dy": PlatformEnum.DOUYIN,
|
||
"ks": PlatformEnum.KUAISHOU,
|
||
"bili": PlatformEnum.BILIBILI,
|
||
"wb": PlatformEnum.WEIBO,
|
||
"tieba": PlatformEnum.TIEBA,
|
||
"zhihu": PlatformEnum.ZHIHU,
|
||
}
|
||
|
||
# Timeout used when the caller does not care; tasks carry their own.
|
||
DEFAULT_RUN_TIMEOUT_SECONDS = 3600
|
||
|
||
|
||
def build_target_url(value: str, kind: str, platform: str) -> str:
|
||
"""Turn a stored target into a URL the crawler's parser accepts.
|
||
|
||
Always emits a full URL rather than a bare id: both platforms' parsers accept
|
||
a bare id only in a narrower form, so the URL is the safer universal input.
|
||
The shape itself is platform-specific and comes from ``adapters``.
|
||
"""
|
||
spec = adapters.adapter(platform)
|
||
path = spec.creator_path if kind == MODE_CREATOR else spec.note_path
|
||
return f"{spec.web_base}{path}/{value}"
|
||
|
||
|
||
def build_target_urls(
|
||
mode: str, targets: Iterable[MonitorTarget], platform: str
|
||
) -> List[str]:
|
||
"""存储的目标 -> 爬虫接受的 URL。
|
||
|
||
``xsec_token`` 只有小红书有,而且是会过期的刷新令牌 —— 有就带上,没有就算了。
|
||
抖音恒为空,所以这一段对它是天然的 no-op,不需要平台分支。
|
||
"""
|
||
urls = []
|
||
for target in targets:
|
||
url = build_target_url(target.external_id, target.kind, platform)
|
||
if target.xsec_token:
|
||
url = f"{url}?xsec_token={target.xsec_token}"
|
||
if target.xsec_source:
|
||
url = f"{url}&xsec_source={target.xsec_source}"
|
||
urls.append(url)
|
||
return urls
|
||
|
||
|
||
async def _strategy_settings(session, platform: str) -> dict:
|
||
"""Crawl-strategy and proxy settings for one platform.
|
||
|
||
Per-platform because the values genuinely differ: what is a safe request
|
||
interval on one site is a rate limit on another. Read per run rather than
|
||
cached, so a change takes effect on the next scheduled run.
|
||
"""
|
||
return {
|
||
"enable_sub_comments": bool(
|
||
await app_settings.get_value(session, "enable_sub_comments", platform, False)
|
||
),
|
||
"crawl_sleep_sec": int(
|
||
await app_settings.get_value(session, "crawl_sleep_sec", platform, 2)
|
||
),
|
||
"enable_ip_proxy": bool(
|
||
await app_settings.get_value(session, "enable_ip_proxy", platform, False)
|
||
),
|
||
"proxy_provider": await app_settings.get_value(
|
||
session, "proxy_provider", platform, "kuaidaili"
|
||
),
|
||
"proxy_pool_count": int(
|
||
await app_settings.get_value(session, "proxy_pool_count", platform, 2)
|
||
),
|
||
"static_proxy_url": await app_settings.get_value(
|
||
session, "static_proxy_url", platform, ""
|
||
),
|
||
}
|
||
|
||
|
||
def _write_cookie_file(path: Path, cookie: str) -> None:
|
||
path.parent.mkdir(parents=True, exist_ok=True)
|
||
path.write_text(cookie, encoding="utf-8")
|
||
|
||
|
||
def _remove_cookie_file(path: Path) -> None:
|
||
"""Best-effort removal; the cookie is a credential, do not leave it around."""
|
||
try:
|
||
os.remove(path)
|
||
except OSError:
|
||
pass
|
||
|
||
|
||
async def execute_task(task_id: int, trigger: str = "manual") -> IngestResult:
|
||
"""Run one monitoring cycle for ``task_id`` and ingest its output.
|
||
|
||
Split into three phases with separate short-lived DB sessions so no
|
||
transaction is held open across the multi-minute subprocess run.
|
||
"""
|
||
# --- Phase 1: book the run and work out where its output goes -------------
|
||
async with get_session() as session:
|
||
task = await session.get(MonitorTask, task_id)
|
||
if task is None:
|
||
raise ValueError(f"Monitor task {task_id} not found")
|
||
|
||
targets = [target for target in task.targets if target.enabled]
|
||
if not targets:
|
||
raise ValueError(f"Monitor task {task_id} has no enabled targets")
|
||
|
||
platform = task.platform
|
||
urls = build_target_urls(task.mode, targets, platform)
|
||
cookie = await get_cookie(session, platform)
|
||
strategy = await _strategy_settings(session, platform)
|
||
# System-wide switch. On a headless server the crawler must attach to the
|
||
# Chrome already listening on the debug port -- that browser is where the
|
||
# operator scanned the login QR, so its profile is the login. Default False
|
||
# keeps desktop runs launching a private browser exactly as before.
|
||
cdp_enabled = await app_settings.get_value(session, "cdp_enabled", fallback=False)
|
||
|
||
run = MonitorRun(
|
||
task_id=task.id,
|
||
trigger=trigger,
|
||
status=RUN_PENDING,
|
||
phase=task.mode,
|
||
save_data_path="",
|
||
queued_at=get_current_timestamp(),
|
||
not_before=0,
|
||
max_comments_count=task.max_comments_count if task.enable_comments else 0,
|
||
)
|
||
session.add(run)
|
||
await session.flush()
|
||
|
||
run_id = run.id
|
||
out_dir = MONITOR_RUNS_DIR / str(task.id) / str(run_id)
|
||
run.save_data_path = str(out_dir)
|
||
|
||
# Snapshot the values the subprocess needs; `task` is detached after commit.
|
||
mode = task.mode
|
||
enable_comments = task.enable_comments
|
||
max_notes_count = task.max_notes_count
|
||
max_comments_count = task.max_comments_count
|
||
timeout_seconds = task.run_timeout_seconds
|
||
# 抖音的作品列表接口被那道真校验挡着(见 douyin_fetch),拿不到列表时就靠这些
|
||
# 已知的 aweme_id 逐条刷新 —— 新作品发现不了,但已有作品的指标还能继续更新。
|
||
known_aweme_ids = list(
|
||
await session.scalars(
|
||
select(MonitorNote.note_id).where(MonitorNote.task_id == task.id)
|
||
)
|
||
)
|
||
|
||
# --- Phase 2: run the collection outside any transaction ------------------
|
||
# 抖音走**进程内 HTTP 客户端**,不起 Playwright 子进程:那边会构造一大串自相矛盾的
|
||
# 浏览器指纹参数(参数说 Mac + Chrome 125、UA 说 Linux + Chrome 155),网关回一个
|
||
# 200 + 空 body,然后被翻译成「account blocked」—— 看着像账号被封,其实什么都不是。
|
||
# 见 douyin_fetch / douyin_api。
|
||
in_process_tail: List[str] = []
|
||
if platform == adapters.PLATFORM_DY:
|
||
fetched = await douyin_fetch.collect(
|
||
out_dir,
|
||
platform=platform,
|
||
mode=mode,
|
||
limit=max_notes_count,
|
||
want_comments=enable_comments,
|
||
comment_limit=max_comments_count,
|
||
targets=targets,
|
||
known_aweme_ids=known_aweme_ids,
|
||
cookie=cookie,
|
||
)
|
||
in_process_tail = list(fetched["errors"])
|
||
# 一条都没采到 = 这一轮失败,并把**真因**当作退出诊断传下去。否则它会掉进
|
||
# ingest 的「疑似登录失效」分支 —— 又骗人一次,正是这套东西一直在犯的毛病。
|
||
exit_code = 1 if (fetched["errors"] and not fetched["notes"]) else 0
|
||
if fetched["errors"] and fetched["notes"]:
|
||
# 有产物但带着错误,说明走了退化路径(比如作品列表被挡,只刷新了已知作品)。
|
||
# 这一轮状态是成功,但**不是**一切正常 —— 得留下痕迹,否则没人知道新作品
|
||
# 其实没在发现。
|
||
print(
|
||
"[monitor.runner] 抖音采集部分失败:"
|
||
+ ";".join(fetched["errors"])[:300]
|
||
)
|
||
else:
|
||
cookie_file = out_dir / ".cookies"
|
||
_write_cookie_file(cookie_file, cookie)
|
||
|
||
request = CrawlerStartRequest(
|
||
platform=_PLATFORM_ENUM[platform],
|
||
login_type=LoginTypeEnum.COOKIE,
|
||
crawler_type=CrawlerTypeEnum.CREATOR if mode == MODE_CREATOR else CrawlerTypeEnum.DETAIL,
|
||
creator_ids=",".join(urls) if mode == MODE_CREATOR else "",
|
||
specified_ids=",".join(urls) if mode != MODE_CREATOR else "",
|
||
start_page=1,
|
||
enable_comments=enable_comments,
|
||
enable_sub_comments=strategy["enable_sub_comments"],
|
||
enable_media=False,
|
||
save_option=SaveDataOptionEnum.JSONL,
|
||
cookies="",
|
||
headless=True,
|
||
max_notes_count=max_notes_count,
|
||
max_comments_count=max_comments_count,
|
||
# Isolate this run's output: the crawler names files by date only, so
|
||
# otherwise same-day runs would append into one shared file.
|
||
save_data_path=str(out_dir),
|
||
# Attach to the browser already running on CDP_DEBUG_PORT when the
|
||
# operator enabled it; otherwise launch a private, throwaway browser.
|
||
enable_cdp_mode=cdp_enabled,
|
||
# Only injecting web_session is not enough to sign requests from a cold
|
||
# browser profile.
|
||
inject_all_cookies=True,
|
||
save_login_state=True,
|
||
cookies_file=str(cookie_file),
|
||
max_concurrency_num=1,
|
||
# Strategy + proxy, surfaced on the Settings page.
|
||
crawler_max_sleep_sec=strategy["crawl_sleep_sec"],
|
||
enable_ip_proxy=strategy["enable_ip_proxy"],
|
||
ip_proxy_pool_count=strategy["proxy_pool_count"],
|
||
ip_proxy_provider_name=strategy["proxy_provider"],
|
||
static_proxy_url=strategy["static_proxy_url"] or None,
|
||
)
|
||
|
||
async with get_session() as session:
|
||
run = await session.get(MonitorRun, run_id)
|
||
if run is not None:
|
||
run.status = RUN_RUNNING
|
||
run.started_at = get_current_timestamp()
|
||
|
||
try:
|
||
exit_code = await crawler_manager.run_and_wait(request, timeout=timeout_seconds)
|
||
finally:
|
||
_remove_cookie_file(cookie_file)
|
||
|
||
# 失败时的诊断来源。抖音那条路没有子进程,尾巴就是它自己报的错 —— **别去读
|
||
# crawler_manager 的尾巴**,那里面是上一轮别的平台留下的东西,会张冠李戴。
|
||
output_tail = (
|
||
in_process_tail
|
||
if platform == adapters.PLATFORM_DY
|
||
else crawler_manager.get_output_tail()
|
||
)
|
||
|
||
# --- Phase 3: ingest ------------------------------------------------------
|
||
async with get_session() as session:
|
||
run = await session.get(MonitorRun, run_id)
|
||
task = await session.get(MonitorTask, task_id)
|
||
if run is None or task is None:
|
||
raise ValueError(f"Run {run_id} or task {task_id} vanished during execution")
|
||
|
||
# 目录名按平台解析 —— 抖音的平台 id 是 dy 而产物目录是 douyin,写死就永远判不准。
|
||
if exit_code == -1 and not (out_dir / adapters.artifact_dir(task.platform)).exists():
|
||
# run_and_wait returns -1 when the process could not start or timed out.
|
||
run.status = RUN_TIMEOUT
|
||
run.finished_at = get_current_timestamp()
|
||
run.exit_code = exit_code
|
||
run.error_message = "Run was killed by timeout or failed to start"
|
||
# -1 同时代表「超时」和「根本没起来」,两者要查的东西完全不同。输出尾巴里
|
||
# 有异常就带上它,否则运行历史里只能看到这句没有信息量的话。
|
||
cause = diagnose_failure(output_tail)
|
||
if cause:
|
||
run.error_message = f"{run.error_message};原因:{cause}"
|
||
result = IngestResult(status=RUN_TIMEOUT, error=run.error_message)
|
||
else:
|
||
run.exit_code = exit_code
|
||
run.finished_at = get_current_timestamp()
|
||
# 把爬虫输出的尾巴交给 ingest:退出码本身说明不了问题,运行历史里要显示的
|
||
# 是真正的报错(比如抖音的 DataFetchError: account blocked)。
|
||
result = await ingest_run(
|
||
session, run, task, out_dir, output_tail=output_tail
|
||
)
|
||
|
||
# A run that authenticated fine is the only useful signal that the
|
||
# stored cookie still works.
|
||
if result.notes_fetched > 0:
|
||
await mark_cookie_ok(session, task.platform)
|
||
|
||
task.last_run_at = run.finished_at
|
||
task.last_status = result.status
|
||
task.last_error = result.error
|
||
|
||
# --- Phase 4: notify ------------------------------------------------------
|
||
# Runs after the ingest transaction has committed, in its own session. A push
|
||
# failure must never roll back collected data, and notify_run() swallows its
|
||
# own errors for the same reason.
|
||
async with get_session() as session:
|
||
task = await session.get(MonitorTask, task_id)
|
||
run = await session.get(MonitorRun, run_id)
|
||
if task is not None and run is not None:
|
||
await notify.notify_run(session, task, run)
|
||
|
||
# --- Phase 5: 封面落盘 ------------------------------------------------------
|
||
# 也放在事务之外。封面地址带签名、会过期(实测隔天即 403),落盘之后才与签名无关。
|
||
# 下载慢且可能失败,占着一个入库事务是不合适的;失败也不影响本轮数据。
|
||
try:
|
||
async with get_session() as session:
|
||
saved = await covers.cache_pending(session, task_id)
|
||
if saved:
|
||
print(f"[monitor.runner] 缓存了 {saved} 张作品封面")
|
||
except Exception as exc: # noqa: BLE001 - 封面拿不到不该让整轮失败
|
||
print(f"[monitor.runner] 封面缓存失败:{exc}")
|
||
|
||
return result
|