| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163 |
- """帖子级「创作知识 / 非创作知识」分类(收紧版)+ 提取知识点。
- 判据已收紧到「**自媒体内容创作**」,并显式排除三类越界(见 prompts):① 应试/学术写作
- ② 制作/工具操作(=制作知识)③ 学科知识/评论/作品本身。两套真实提示词:
- 图文(小红书/微信):prompts/classify_imgtext.txt(标题+正文+全部图喂 Gemini)
- 视频(抖音):prompts/classify_video.txt(整段视频原生喂 Gemini,大视频先 ffmpeg 压 480p 保音轨)
- 都输出 {is_empty, reason, knowledge}:is_empty 即创作闸;is_empty=false 时连带提炼出具体创作知识点。
- 结果按 url 写 app.db 的 post_class(upsert 覆盖,可重判)。并发;与微信补下载并行安全(busy_timeout)。
- 只读 prompts、走 OpenRouter——不改 skill/creation_knowledge,不解构、不入 ingest。
- 用法:PYTHONPATH=. CK_ENV_FILE=.env python -m acquisition.classify [平台名...]
- """
- from __future__ import annotations
- import base64
- import concurrent.futures as cf
- import json
- import subprocess
- import sys
- import tempfile
- import time
- from pathlib import Path
- import httpx
- from acquisition import store
- from core.config import Settings
- from core.prompts import load_prompt
- ROOT = Path(__file__).resolve().parent.parent
- PLATFORMS = ["xiaohongshu", "weixin", "douyin"] # 默认重判全部;可传平台名覆盖
- IMG_WORKERS = 8 # 图文并发
- VID_WORKERS = 3 # 视频并发(含 ffmpeg 压制,别太高)
- MAX_CARDS = 12 # 图文最多送几张图
- COMPRESS_OVER_MB = 12 # 视频超过此大小先压再喂
- COMPRESS_H = 480
- def _data_url(public_path: str):
- fs = ROOT / public_path.lstrip("/")
- if not fs.exists():
- return None
- return "data:image/jpeg;base64," + base64.b64encode(fs.read_bytes()).decode()
- def _compress(mp4: Path) -> Path:
- """大视频压到 480p(保留口播音轨)→ 临时 mp4;失败回原文件。"""
- out = Path(tempfile.gettempdir()) / f"ck_{mp4.parent.name}_480.mp4"
- try:
- import imageio_ffmpeg
- ff = imageio_ffmpeg.get_ffmpeg_exe()
- subprocess.run([ff, "-y", "-i", str(mp4), "-vf", f"scale=-2:{COMPRESS_H}",
- "-c:v", "libx264", "-crf", "30", "-preset", "veryfast",
- "-c:a", "aac", "-b:a", "64k", str(out)],
- capture_output=True, timeout=300)
- if out.exists() and out.stat().st_size > 0:
- return out
- except Exception:
- pass
- return mp4
- def _judge(messages: list, settings: Settings, timeout: float) -> tuple:
- """调 Gemini(OpenRouter,强制 JSON)解析 {is_empty, reason, knowledge},带重试。
- 返回 (is_creation 1/0/None, reason, knowledge, points)。"""
- api = settings.openrouter_base_url.rstrip("/") + "/chat/completions"
- headers = {"Authorization": f"Bearer {settings.openrouter_api_key}", "Content-Type": "application/json"}
- payload = {"model": settings.video_model, "messages": messages,
- "response_format": {"type": "json_object"}}
- last = ""
- for attempt in range(3):
- try:
- resp = httpx.post(api, headers=headers, json=payload, timeout=timeout)
- if resp.status_code == 200:
- d = json.loads(resp.json()["choices"][0]["message"]["content"])
- if bool(d.get("is_empty")):
- return 0, str(d.get("reason", ""))[:60], "", ""
- return 1, str(d.get("reason", ""))[:60], str(d.get("knowledge", "") or ""), ""
- last = f"http {resp.status_code}"
- except Exception as exc:
- last = str(exc)[:50]
- time.sleep(2 * (attempt + 1))
- return None, f"判定失败: {last}", "", ""
- def classify_imgtext(p: dict, settings: Settings) -> tuple:
- """图文:标题+正文+全部图,用收紧的 classify_imgtext.txt 判 is_empty 并提取知识点。"""
- user = [{"type": "text", "text": f"平台:{p.get('platform')}\n标题:{p.get('title', '')}\n"
- f"正文:{(p.get('body_text') or '')[:1500]}\n(下附帖子图片,请一并看完)"}]
- for im in (p.get("images") or [])[:MAX_CARDS]:
- u = _data_url(im)
- if u:
- user.append({"type": "image_url", "image_url": {"url": u}})
- messages = [{"role": "system", "content": load_prompt("classify_imgtext")},
- {"role": "user", "content": user}]
- return _judge(messages, settings, timeout=120)
- def classify_video(p: dict, settings: Settings) -> tuple:
- """视频:看完整段视频,用收紧的 classify_video.txt 判 is_empty 并提取知识点。"""
- rel = p.get("video") or ""
- mp4 = ROOT / rel.lstrip("/")
- if not rel or not mp4.exists():
- return None, "无本地视频", "", ""
- use = _compress(mp4) if mp4.stat().st_size > COMPRESS_OVER_MB * 1048576 else mp4
- try:
- media = "data:video/mp4;base64," + base64.b64encode(use.read_bytes()).decode()
- finally:
- if use != mp4:
- try:
- use.unlink()
- except Exception:
- pass
- messages = [{"role": "system", "content": load_prompt("classify_video")},
- {"role": "user", "content": [{"type": "text", "text": "判断这条视频是不是创作知识。"},
- {"type": "video_url", "video_url": {"url": media}}]}]
- return _judge(messages, settings, timeout=300)
- def _safe(fn, *a) -> tuple:
- try:
- return fn(*a)
- except Exception as exc:
- return None, f"判定失败: {str(exc)[:60]}", "", ""
- def main() -> None:
- settings = Settings.from_env()
- conn = store.connect()
- platforms = sys.argv[1:] or PLATFORMS
- posts = store.posts_to_classify(conn, platforms)
- imgs = [p for p in posts if p["platform"] != "douyin"]
- vids = [p for p in posts if p["platform"] == "douyin"]
- total = len(posts)
- print(f"收紧重判:图文 {len(imgs)}(并发{IMG_WORKERS})+ 抖音视频 {len(vids)}(并发{VID_WORKERS},大视频先压)")
- ts = int(time.time())
- done = {"n": 0, "fail": 0}
- def _write(p, res):
- ic, reason, knowledge, points = res
- if ic is None:
- done["fail"] += 1
- else:
- store.upsert_class(conn, p["url"], ic, reason, ts, knowledge, points)
- done["n"] += 1
- if done["n"] % 20 == 0:
- print(f" {done['n']}/{total}(失败 {done['fail']})")
- with cf.ThreadPoolExecutor(IMG_WORKERS) as ex:
- futs = {ex.submit(_safe, classify_imgtext, p, settings): p for p in imgs}
- for fut in cf.as_completed(futs):
- _write(futs[fut], fut.result())
- with cf.ThreadPoolExecutor(VID_WORKERS) as ex:
- futs = {ex.submit(_safe, classify_video, p, settings): p for p in vids}
- for fut in cf.as_completed(futs):
- _write(futs[fut], fut.result())
- c = store.class_counts(conn)
- conn.close()
- print(f"完成:创作知识 {c['creation']} / 非创作知识 {c['non_creation']}(本轮失败 {done['fail']},可重跑补判)")
- if __name__ == "__main__":
- main()
|