diff --git a/README.md b/README.md index 0c11c4553..f31032d9c 100644 --- a/README.md +++ b/README.md @@ -296,20 +296,48 @@ MediaCrawler 支持多种数据存储方式,包括 CSV、JSON、JSONL、Excel ## 💰 赞助商展示 - - -
-TikHub.io 提供 900+ 高稳定性数据接口,覆盖 TK、DY、XHS、Y2B、Ins、X 等 14+ 海内外主流平台,支持用户、内容、商品、评论等多维度公开数据 API,并配套 4000 万+ 已清洗结构化数据集,使用邀请码 cfzyejV9 注册并充值,即可额外获得 $2 赠送额度。 -
-
-
- - -Atlas Cloud -Atlas Cloud - -
-Atlas Cloud 是一个全模态 AI 推理平台,让开发者通过统一的 AI API 访问视频生成、图像生成和 LLM API,无需分别维护多个厂商集成,即可调用 300+ 精选模型。Atlas Cloud 最新推出 coding plan 优惠,为开发者提供更具性价比的 API 访问预算。 + + + + + + + + + + + + + + + + + + + + + + + + + +
赞助商介绍
+ TikHub + + TikHub.io 提供 900+ 高稳定性数据接口,覆盖 TK、DY、XHS、Y2B、Ins、X 等 14+ 海内外主流平台,支持用户、内容、商品、评论等多维度公开数据 API,并配套 4000 万+ 已清洗结构化数据集,使用邀请码 cfzyejV9 注册并充值,即可额外获得 $2 赠送额度。 +
+ Atlas CloudAtlas Cloud + + Atlas Cloud 是一个全模态 AI 推理平台,让开发者通过统一的 AI API 访问视频生成、图像生成和 LLM API,无需分别维护多个厂商集成,即可调用 300+ 精选模型。Atlas Cloud 最新推出 coding plan 优惠,为开发者提供更具性价比的 API 访问预算。 +
+ Bloome + + Bloome 是一个 AI Agent IM 平台——让多个 AI agent(Claude、ChatGPT、DeepSeek 等)和你在同一个对话里像团队成员一样协作,自动分工、互相校对,直接生成表格、文档与可视化看板。零配置、云端运行,网页和手机都能用,还能把配好的 agent 一键分享给团队。👉 试试 Bloome +
+ NodeMaven + + NodeMaven 提供稳定可靠的高质量代理服务,适用于自动化、网页抓取、SEO 研究和社交媒体管理。服务支持 99.9% 可用性、最长 7 天的粘性会话、IP 质量筛选(所有代理的欺诈评分均低于 97%)、无需 KYC,以及最高 10% 的流量返现。MediaCrawler 用户使用优惠码 CRAWLER35 可享移动和住宅代理 35% 折扣,使用 CRAWLER40 可享 ISP(静态)代理 40% 折扣。👉 访问 NodeMaven +
--- diff --git a/README_en.md b/README_en.md index 36a2d4874..fb0d96244 100644 --- a/README_en.md +++ b/README_en.md @@ -278,20 +278,48 @@ MediaCrawler supports multiple data storage methods, including CSV, JSON, JSONL, ### 💰 Sponsor Display - - -
-TikHub.io provides 900+ highly stable data interfaces, covering 14+ mainstream domestic and international platforms including TK, DY, XHS, Y2B, Ins, X, etc. Supports multi-dimensional public data APIs for users, content, products, comments, etc., with 40M+ cleaned structured datasets. Use invitation code cfzyejV9 to register and recharge, and get an additional $2 bonus. -
-
-
- - -Atlas Cloud -Atlas Cloud - -
-Atlas Cloud is a full-modal AI inference platform that gives developers a single AI API to access video generation, image generation, and LLM APIs. Instead of managing multiple vendor integrations, you connect once and get unified access to 300+ curated models across all modalities. Check out Atlas Cloud's new coding plan promotion for more budget-friendly API access. + + + + + + + + + + + + + + + + + + + + + + + + + +
SponsorIntroduction
+ TikHub + + TikHub.io provides 900+ highly stable data interfaces, covering 14+ mainstream domestic and international platforms including TK, DY, XHS, Y2B, Ins, X, etc. Supports multi-dimensional public data APIs for users, content, products, comments, etc., with 40M+ cleaned structured datasets. Use invitation code cfzyejV9 to register and recharge, and get an additional $2 bonus. +
+ Atlas CloudAtlas Cloud + + Atlas Cloud is a full-modal AI inference platform that gives developers a single AI API to access video generation, image generation, and LLM APIs. Instead of managing multiple vendor integrations, you connect once and get unified access to 300+ curated models across all modalities. Check out Atlas Cloud's new coding plan promotion for more budget-friendly API access. +
+ Bloome + + Bloome is an AI Agent IM platform — multiple AI agents (Claude, ChatGPT, DeepSeek, etc.) collaborate with you in a single conversation like team members, automatically dividing up the work and cross-checking each other, and directly producing tables, documents, and visual dashboards. Zero config, runs in the cloud, works on both web and mobile, and you can share your configured agents with your team in one click. 👉 Try Bloome +
+ NodeMaven + + NodeMaven provides reliable, high-quality proxies for automation, web scraping, SEO research, and social media management. Features include 99.9% uptime, sticky sessions up to 7 days, IP filtering across all proxies (fraud score below 97%), no KYC, and traffic cashback of up to 10%. MediaCrawler users get 35% off mobile and residential proxies with code CRAWLER35, and 40% off ISP (static) proxies with code CRAWLER40. 👉 Visit NodeMaven +
--- diff --git a/README_es.md b/README_es.md index 66974c839..f971f1faa 100644 --- a/README_es.md +++ b/README_es.md @@ -278,11 +278,40 @@ MediaCrawler soporta múltiples métodos de almacenamiento de datos, incluyendo ### 💰 Exhibición de Patrocinadores - - -
-TikHub.io proporciona 900+ interfaces de datos altamente estables, cubriendo 14+ plataformas principales nacionales e internacionales incluyendo TK, DY, XHS, Y2B, Ins, X, etc. Soporta APIs de datos públicos multidimensionales para usuarios, contenido, productos, comentarios, etc., con 40M+ conjuntos de datos estructurados limpios. Use el código de invitación cfzyejV9 para registrarse y recargar, y obtenga $2 adicionales de bonificación. -
+ + + + + + + + + + + + + + + + + + + + + +
PatrocinadorIntroducción
+ TikHub + + TikHub.io proporciona 900+ interfaces de datos altamente estables, cubriendo 14+ plataformas principales nacionales e internacionales incluyendo TK, DY, XHS, Y2B, Ins, X, etc. Soporta APIs de datos públicos multidimensionales para usuarios, contenido, productos, comentarios, etc., con 40M+ conjuntos de datos estructurados limpios. Use el código de invitación cfzyejV9 para registrarse y recargar, y obtenga $2 adicionales de bonificación. +
+ Bloome + + Bloome es una plataforma de IM de agentes de IA: varios agentes de IA (Claude, ChatGPT, DeepSeek, etc.) colaboran contigo en una misma conversación como miembros de un equipo, dividiéndose el trabajo automáticamente y revisándose entre sí, y generando directamente tablas, documentos y paneles visuales. Sin configuración, funciona en la nube, disponible tanto en web como en móvil, y puedes compartir tus agentes configurados con tu equipo con un solo clic. 👉 Prueba Bloome +
+ NodeMaven + + NodeMaven ofrece proxies fiables y de alta calidad para automatización, web scraping, investigación SEO y gestión de redes sociales. El servicio incluye una disponibilidad del 99,9%, sesiones persistentes de hasta 7 días, filtrado de IP en todos los proxies (puntuación de fraude inferior al 97%), sin KYC y reembolso de hasta el 10% del tráfico. Los usuarios de MediaCrawler obtienen un 35% de descuento en proxies móviles y residenciales con el código CRAWLER35, y un 40% de descuento en proxies ISP (estáticos) con CRAWLER40. 👉 Visita NodeMaven +
--- diff --git a/api/main.py b/api/main.py index af49b152f..2fd0939e3 100644 --- a/api/main.py +++ b/api/main.py @@ -50,10 +50,12 @@ app.add_middleware( CORSMiddleware, allow_origins=[ - "http://localhost:5173", # Vite dev server + "http://localhost:5173", # Vite dev server (original) "http://localhost:3000", # Backup port "http://127.0.0.1:5173", "http://127.0.0.1:3000", + "http://localhost:15173", # Vite dev server (custom port, avoid conflict) + "http://127.0.0.1:15173", ], allow_credentials=True, allow_methods=["*"], @@ -202,4 +204,4 @@ async def get_config_options(): if __name__ == "__main__": - uvicorn.run(app, host="0.0.0.0", port=8080) + uvicorn.run(app, host="0.0.0.0", port=18080) diff --git a/api/routers/crawler.py b/api/routers/crawler.py index eead9e1ae..f429c81a0 100644 --- a/api/routers/crawler.py +++ b/api/routers/crawler.py @@ -16,13 +16,32 @@ # 详细许可条款请参阅项目根目录下的LICENSE文件。 # 使用本代码即表示您同意遵守上述原则和LICENSE中的所有条款。 +import os +from pathlib import Path + from fastapi import APIRouter, HTTPException -from ..schemas import CrawlerStartRequest, CrawlerStatusResponse +from ..schemas import CrawlerStartRequest, CrawlerStatusResponse, ClearHistoryRequest from ..services import crawler_manager +from ..services import run_history router = APIRouter(prefix="/crawler", tags=["crawler"]) +# 项目根目录 / data 目录(清除数据文件用) +PROJECT_ROOT = Path(__file__).parent.parent.parent +DATA_DIR = PROJECT_ROOT / "data" + +# 平台短名 → 存储目录名 / ORM 表名前缀 映射 +PLATFORM_MAP = { + "dy": ("douyin", "douyin"), + "xhs": ("xhs", "xhs"), + "ks": ("kuaishou", "kuaishou"), + "bili": ("bilibili", "bilibili"), + "wb": ("weibo", "weibo"), + "tieba": ("tieba", "tieba"), + "zhihu": ("zhihu", "zhihu"), +} + @router.post("/start") async def start_crawler(request: CrawlerStartRequest): @@ -61,3 +80,132 @@ async def get_logs(limit: int = 100): """Get recent logs""" logs = crawler_manager.logs[-limit:] if limit > 0 else crawler_manager.logs return {"logs": [log.model_dump() for log in logs]} + + +@router.get("/history") +async def get_run_history(limit: int = 50): + """获取爬取运行历史,最新在前,上限 limit""" + runs = run_history.get_recent_runs(limit=limit) + return {"runs": runs} + + +@router.delete("/history") +async def clear_history(request: ClearHistoryRequest): + """ + 清除历史数据。可分别清:数据文件 / DB 表 / 运行清单。 + + 爬取进行中拒绝清除(409),避免破坏追加模式文件。 + """ + # 爬取中拒绝 + if crawler_manager.status == "running": + raise HTTPException(status_code=409, detail="Cannot clear history while crawler is running") + + result = {"deleted_files": 0, "truncated_tables": [], "cleared_runs": False} + + # 1. 清数据文件 + if request.clear_files: + result["deleted_files"] = _delete_data_files(request.platform) + + # 2. 清 DB 表 + if request.clear_db: + result["truncated_tables"] = _truncate_db_tables(request.platform) + + # 3. 清运行清单 + if request.clear_runs: + run_history.clear_runs() + result["cleared_runs"] = True + + return result + + +def _delete_data_files(platform: str | None) -> int: + """删除 data/ 下的数据文件,可选按平台过滤。跳过 .run_history.json 清单本身。""" + if not DATA_DIR.exists(): + return 0 + supported_extensions = {".json", ".jsonl", ".csv", ".xlsx", ".xls"} + # 平台短名 → 存储目录名 + platform_dir = PLATFORM_MAP.get(platform, ("", ""))[0] if platform else "" + deleted = 0 + for root, _dirs, filenames in os.walk(DATA_DIR): + root_path = Path(root) + try: + rel = str(root_path.relative_to(DATA_DIR)) + except ValueError: + continue + # 平台过滤:目录路径需包含平台存储名 + if platform_dir and platform_dir not in rel.lower(): + continue + for filename in filenames: + if filename == ".run_history.json": + continue # 跳过运行清单本身 + if filename == ".bgm_tags.json": + continue # 跳过 BGM 场景标签标注文件 + file_path = root_path / filename + if file_path.suffix.lower() not in supported_extensions: + continue + try: + file_path.unlink() + deleted += 1 + except Exception: + continue + return deleted + + +def _truncate_db_tables(platform: str | None) -> list[str]: + """truncate ORM 表,可选按平台过滤。文件存储模式(无 engine)返回空列表。""" + try: + from database.db_session import get_async_engine + from database import models + import asyncio + from sqlalchemy import text + except Exception: + return [] + + engine = get_async_engine() + if engine is None: + # 文件存储模式(jsonl/json/csv/excel),无 DB,跳过 + return [] + + # 收集所有 ORM 表名 + all_tables = [t.name for t in models.Base.metadata.tables.values()] + # 平台过滤:表名前缀匹配 + if platform: + prefix = PLATFORM_MAP.get(platform, ("", ""))[1] + if not prefix: + return [] + target_tables = [t for t in all_tables if t.startswith(prefix)] + else: + target_tables = all_tables + + truncated = [] + + async def _do_truncate(): + async with engine.begin() as conn: + for table_name in target_tables: + # 表名来自代码常量,无注入风险,但仍用 text 绑定 + await conn.execute(text(f'DELETE FROM "{table_name}"')) + truncated.append(table_name) + + try: + asyncio.get_event_loop().run_until_complete(_do_truncate()) + except RuntimeError: + # 已有事件循环时(FastAPI 上下文),用新线程跑 + import threading + result_holder = {"err": None} + + def _run(): + try: + new_loop = asyncio.new_event_loop() + asyncio.set_event_loop(new_loop) + new_loop.run_until_complete(_do_truncate()) + new_loop.close() + except Exception as e: + result_holder["err"] = e + + t = threading.Thread(target=_run) + t.start() + t.join() + if result_holder["err"]: + raise result_holder["err"] + + return truncated diff --git a/api/routers/data.py b/api/routers/data.py index 7dc81aff6..7d540f384 100644 --- a/api/routers/data.py +++ b/api/routers/data.py @@ -17,12 +17,14 @@ # 使用本代码即表示您同意遵守上述原则和LICENSE中的所有条款。 import os +import re +import glob import json from pathlib import Path from typing import Optional -from fastapi import APIRouter, HTTPException -from fastapi.responses import FileResponse +from fastapi import APIRouter, HTTPException, Request +from fastapi.responses import FileResponse, StreamingResponse router = APIRouter(prefix="/data", tags=["data"]) @@ -37,12 +39,16 @@ def get_file_info(file_path: Path) -> dict: # Try to get record count try: - if file_path.suffix == ".json": + if file_path.suffix.lower() == ".json": with open(file_path, "r", encoding="utf-8") as f: data = json.load(f) if isinstance(data, list): record_count = len(data) - elif file_path.suffix == ".csv": + elif file_path.suffix.lower() == ".jsonl": + # jsonl 每行一个 JSON 对象,统计非空行数 + with open(file_path, "r", encoding="utf-8") as f: + record_count = sum(1 for line in f if line.strip()) + elif file_path.suffix.lower() == ".csv": with open(file_path, "r", encoding="utf-8") as f: record_count = sum(1 for _ in f) - 1 # Subtract header row except Exception: @@ -65,7 +71,7 @@ async def list_data_files(platform: Optional[str] = None, file_type: Optional[st return {"files": []} files = [] - supported_extensions = {".json", ".csv", ".xlsx", ".xls"} + supported_extensions = {".json", ".jsonl", ".csv", ".xlsx", ".xls"} for root, dirs, filenames in os.walk(DATA_DIR): root_path = Path(root) @@ -115,12 +121,28 @@ async def get_file_content(file_path: str, preview: bool = True, limit: int = 10 if preview: # Return preview data try: - if full_path.suffix == ".json": + if full_path.suffix.lower() == ".json": with open(full_path, "r", encoding="utf-8") as f: data = json.load(f) if isinstance(data, list): return {"data": data[:limit], "total": len(data)} return {"data": data, "total": 1} + elif full_path.suffix.lower() == ".jsonl": + # jsonl: 每行一个 JSON 对象,逐行解析,返回前 limit 条 + rows = [] + total = 0 + with open(full_path, "r", encoding="utf-8") as f: + for line in f: + line = line.strip() + if not line: + continue + total += 1 + if len(rows) < limit: + try: + rows.append(json.loads(line)) + except json.JSONDecodeError: + continue + return {"data": rows, "total": total} elif full_path.suffix == ".csv": import csv with open(full_path, "r", encoding="utf-8") as f: @@ -200,7 +222,7 @@ async def get_data_stats(): "by_type": {} } - supported_extensions = {".json", ".csv", ".xlsx", ".xls"} + supported_extensions = {".json", ".jsonl", ".csv", ".xlsx", ".xls"} for root, dirs, filenames in os.walk(DATA_DIR): root_path = Path(root) @@ -228,3 +250,411 @@ async def get_data_stats(): continue return stats + + +# ==================== BGM 播放相关 ==================== + +# aweme_id 合法字符(防路径穿越) +_AWEME_ID_RE = re.compile(r"^[A-Za-z0-9]+$") + +# BGM 文件扩展名优先级(glob bgm.* 后按此排序取第一个) +_BGM_EXT_PRIORITY = [".m4a", ".mp3", ".mp4", ".m4v"] + +# 扩展名 → MIME 映射 +_BGM_MIME = { + ".m4a": "audio/mp4", + ".mp3": "audio/mpeg", + ".mp4": "audio/mp4", + ".m4v": "audio/mp4", +} + + +@router.get("/bgm/playlist") +async def get_bgm_playlist(run_id: Optional[str] = None): + """ + 读取 BGM 播放清单,按 run 分组返回。 + + 扫描所有 data/douyin/jsonl/search_bgm_playlist_*.jsonl(不只最新),按 aweme_id + 去重(保留最大 add_ts)。每条 track 带 run_id/keyword/add_ts/has_local。 + groups 按 run_id 分组并合并 run_history 元信息(关键词/开始时间/状态)。 + 可用 ?run_id= 过滤单个运行。历史数据无 run_id 归 "unattributed" 组。 + """ + playlist_dir = DATA_DIR / "douyin" / "jsonl" + if not playlist_dir.exists(): + return {"tracks": [], "groups": []} + + # 收集所有 bgm_playlist jsonl,按 mtime 倒序(新文件优先,去重时保留新记录) + candidates = sorted( + playlist_dir.glob("search_bgm_playlist_*.jsonl"), + key=lambda p: p.stat().st_mtime, + reverse=True, + ) + + tracks_by_id: dict[str, dict] = {} + for fpath in candidates: + try: + with open(fpath, "r", encoding="utf-8") as f: + for line in f: + line = line.strip() + if not line: + continue + try: + item = json.loads(line) + except json.JSONDecodeError: + continue + aweme_id = item.get("aweme_id", "") + if not aweme_id: + continue + existing = tracks_by_id.get(aweme_id) + if existing and existing.get("add_ts", 0) >= item.get("add_ts", 0): + continue + tracks_by_id[aweme_id] = item + except Exception: + continue + + # 组装扁平 tracks(带 run_id)+ 按 run 分组 + tracks = [] + groups_map: dict[str, dict] = {} + for aweme_id, item in tracks_by_id.items(): + rid = item.get("run_id") or "unattributed" + if run_id and rid != run_id: + continue + track = { + "aweme_id": aweme_id, + "music_title": item.get("music_title", ""), + "music_author": item.get("music_author", ""), + "music_duration": item.get("music_duration", 0), + "aweme_url": item.get("aweme_url", ""), + "has_local": _resolve_bgm_file(aweme_id) is not None, + "run_id": rid, + "keyword": item.get("keyword", ""), + "add_ts": item.get("add_ts", 0), + } + tracks.append(track) + g = groups_map.setdefault(rid, {"run_id": rid, "keyword": track["keyword"], "tracks": []}) + g["tracks"].append(track) + + groups = _enrich_run_groups(groups_map) + return {"tracks": tracks, "groups": groups} + + +def _resolve_bgm_file(aweme_id: str) -> Optional[Path]: + """ + 按 aweme_id 解析磁盘上的 BGM 音频文件。 + + 不信任 jsonl 里的 local_path(实测可能为空或扩展名不一致), + 直接 glob data/douyin/bgm//bgm.*,按优先级取第一个。 + 返回 None 表示找不到。 + """ + if not _AWEME_ID_RE.match(aweme_id): + return None + bgm_dir = DATA_DIR / "douyin" / "bgm" / aweme_id + if not bgm_dir.exists(): + return None + # glob bgm.* + candidates = list(bgm_dir.glob("bgm.*")) + if not candidates: + return None + # 按优先级排序 + candidates.sort(key=lambda p: _BGM_EXT_PRIORITY.index(p.suffix.lower()) if p.suffix.lower() in _BGM_EXT_PRIORITY else 99) + return candidates[0] + + +# BGM 场景标签标注文件(aweme_id -> 场景名) +_BGM_TAGS_PATH = DATA_DIR / ".bgm_tags.json" + + +def _load_bgm_tags() -> dict: + """读取 BGM 场景标签映射。文件不存在或损坏返回空 dict。""" + if not _BGM_TAGS_PATH.exists(): + return {} + try: + with open(_BGM_TAGS_PATH, "r", encoding="utf-8") as f: + tags = json.load(f) + return tags if isinstance(tags, dict) else {} + except (json.JSONDecodeError, OSError): + return {} + + +def _save_bgm_tags(tags: dict) -> None: + """原子写 BGM 场景标签(tmp + os.replace)。""" + DATA_DIR.mkdir(parents=True, exist_ok=True) + tmp = _BGM_TAGS_PATH.with_suffix(".tmp") + try: + with open(tmp, "w", encoding="utf-8") as f: + json.dump(tags, f, ensure_ascii=False, indent=2) + os.replace(tmp, _BGM_TAGS_PATH) + except Exception: + try: + tmp.unlink(missing_ok=True) + except Exception: + pass + + +@router.get("/bgm/tags") +async def get_bgm_tags(): + """读取全部 BGM 场景标签({aweme_id: scene})。必须在 /bgm/{aweme_id} 之前声明。""" + return {"tags": _load_bgm_tags()} + + +@router.put("/bgm/scene/{aweme_id}") +async def update_bgm_scene(aweme_id: str, body: Optional[dict] = None): + """ + 设置/更新某 BGM 的场景标签。body: {"scene": "婚礼"};空字符串=清除标签。 + 路径用 /bgm/scene/{aweme_id} 避免被 GET /bgm/{aweme_id} 吞掉。 + """ + if not _AWEME_ID_RE.match(aweme_id): + raise HTTPException(status_code=400, detail="Invalid aweme_id") + scene = "" + if body and isinstance(body, dict): + scene = str(body.get("scene", "")).strip() + tags = _load_bgm_tags() + if scene: + tags[aweme_id] = scene + else: + tags.pop(aweme_id, None) + _save_bgm_tags(tags) + return {"aweme_id": aweme_id, "scene": scene} + + +@router.get("/bgm/{aweme_id}") +async def stream_bgm(aweme_id: str, request: Request): + """ + 流式播放指定 aweme_id 的 BGM 音频。 + + 手动实现 HTTP Range 支持(Starlette 0.37 FileResponse 不处理 Range), + 让