mirror of
https://gitverse.ru/kpa39l/vesti.git
synced 2026-09-29 09:55:03 +00:00
feat(crawler): rss-crawler + очередь кандидатов + воркер (фазы 1-3); веб-статус /crawlers; openspec rss-crawler, tg-crawler-publisher-prototype
This commit is contained in:
@@ -0,0 +1,170 @@
|
||||
# VESTI — воркер (фаза 3): параллельная обработка очереди кандидатов.
|
||||
# ThreadPoolExecutor(N): каждый воркер атомарно забирает пост из очереди
|
||||
# (posts WHERE classified=0), классифицирует (keywords → Ollama), пишет
|
||||
# direction/relevance/interest/summary, ставит classified=1. При сбое — возвращает
|
||||
# в очередь (classified=0). Очередь — сама таблица posts, отдельная таблица не нужна.
|
||||
#
|
||||
# Запуск: python -m crawler.worker --workers 8 --limit 200
|
||||
import argparse
|
||||
import os
|
||||
import sqlite3
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
from concurrent.futures import ThreadPoolExecutor, as_completed
|
||||
from pathlib import Path
|
||||
|
||||
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
|
||||
|
||||
from config import DB_PATH, ensure_dirs
|
||||
from classifier.classify import classify_text, call_ollama # noqa: F401 (переиспользуем)
|
||||
|
||||
# Лимит одновременных LLM-запросов: локальная Ollama (qwen3:8b на CPU/GPU) держит
|
||||
# один слот — параллельные 8 запросов = каждый ждёт >CLASSIFY_TIMEOUT и падает.
|
||||
# 8 воркеров = 8 потоков, но к LLM пускаем максимум N (обычно 2) одновременных.
|
||||
LLM_SEM = threading.BoundedSemaphore(int(os.getenv("MAX_CONCURRENT_LLM", "2")))
|
||||
|
||||
|
||||
def claim_post(conn: sqlite3.Connection) -> sqlite3.Row | None:
|
||||
"""Атомарный claim поста из очереди.
|
||||
|
||||
BEGIN IMMEDIATE блокирует других воркеров на время claim; SELECT→UPDATE
|
||||
гарантирует, что строку заберёт ровно один воркер (classified: 0 → -1).
|
||||
Возвращает строку поста или None (очередь пуста).
|
||||
"""
|
||||
conn.execute("BEGIN IMMEDIATE")
|
||||
try:
|
||||
row = conn.execute(
|
||||
"SELECT p.id, p.text, p.views, p.reactions_total, p.is_own, s.slug AS source_slug "
|
||||
"FROM posts p LEFT JOIN sources s ON s.id = p.source_id "
|
||||
"WHERE p.classified IS NULL OR p.classified=0 "
|
||||
"ORDER BY p.is_own DESC, p.id LIMIT 1"
|
||||
).fetchone()
|
||||
if row:
|
||||
conn.execute("UPDATE posts SET classified=-1 WHERE id=?", (row["id"],))
|
||||
conn.commit()
|
||||
return row
|
||||
except Exception:
|
||||
conn.rollback()
|
||||
raise
|
||||
|
||||
|
||||
def process_post(conn: sqlite3.Connection, post_id: int, text: str,
|
||||
views: int, reactions: int, is_own: int, source_slug: str | None = None) -> dict:
|
||||
"""Классифицирует пост и пишет результат. Возвращает result classify_text.
|
||||
|
||||
При исключении — пост возвращается в очередь (classified=0), НЕ теряется.
|
||||
Мультинаправления из словаря (find_directions_all) — в classifications.
|
||||
source_slug (правило SOURCE_RULES, напр. lwn→tech) — направление принудительное.
|
||||
"""
|
||||
res = classify_text({"text": text, "views": views,
|
||||
"reactions_total": reactions, "is_own": is_own,
|
||||
"source_slug": source_slug})
|
||||
conn.execute(
|
||||
"UPDATE posts SET direction=?, relevance=?, interest=?, summary=?, classified=? WHERE id=?",
|
||||
(res["direction"], res["relevance"], res["interest"], res["summary"],
|
||||
1 if res["classified"] else 0, post_id),
|
||||
)
|
||||
# мультинаправления: все словарные попадания → classifications (для fan-out)
|
||||
from classifier.keywords import find_directions_all
|
||||
dirs_all = find_directions_all(text or "", source_slug)
|
||||
if res["direction"] and res["direction"] not in dirs_all:
|
||||
dirs_all.append(res["direction"])
|
||||
for d in dirs_all:
|
||||
exists = conn.execute(
|
||||
"SELECT 1 FROM classifications WHERE post_id=? AND direction=? LIMIT 1",
|
||||
(post_id, d),
|
||||
).fetchone()
|
||||
if not exists:
|
||||
conn.execute(
|
||||
"""INSERT INTO classifications (post_id, direction, relevance, interest, summary, model)
|
||||
VALUES (?,?,?,?,?,?)""",
|
||||
(post_id, d, res["relevance"], res["interest"], res["summary"],
|
||||
res.get("method") or "classify"),
|
||||
)
|
||||
conn.commit()
|
||||
return res
|
||||
|
||||
|
||||
def work_loop(worker_id: int, stats: dict, stop: threading.Event, limit: int):
|
||||
"""Цикл воркера: claim → классификация → запись, пока очередь не пуста/лимит.
|
||||
|
||||
Соединение создаётся ВНУТРИ потока (SQLite: соединение привязано к thread).
|
||||
"""
|
||||
conn = _new_conn()
|
||||
processed = 0
|
||||
try:
|
||||
while not stop.is_set() and processed < limit:
|
||||
row = claim_post(conn)
|
||||
if row is None:
|
||||
break
|
||||
try:
|
||||
# LLM-вызовы ограничены семифором (не перегружаем Ollama)
|
||||
with LLM_SEM:
|
||||
res = process_post(conn, row["id"], row["text"], row["views"],
|
||||
row["reactions_total"], row["is_own"],
|
||||
row["source_slug"])
|
||||
with stats["lock"]:
|
||||
stats["processed"] += 1
|
||||
if res["classified"]:
|
||||
stats["classified"] += 1
|
||||
if res.get("method") == "llm":
|
||||
stats["llm_ok"] += 1
|
||||
except Exception as e:
|
||||
# вернуть в очередь — пост не теряется; зафиксировать счётчик ошибок
|
||||
try:
|
||||
conn.execute("UPDATE posts SET classified=0 WHERE id=?", (row["id"],))
|
||||
conn.commit()
|
||||
except Exception:
|
||||
pass
|
||||
with stats["lock"]:
|
||||
stats["errors"] += 1
|
||||
print(f"[worker {worker_id}] error post {row['id']}: {e}", flush=True)
|
||||
processed += 1
|
||||
finally:
|
||||
conn.close()
|
||||
return processed
|
||||
|
||||
|
||||
def _new_conn() -> sqlite3.Connection:
|
||||
"""Отдельное соединение для воркера (WAL — читатели/писатели параллельны)."""
|
||||
ensure_dirs()
|
||||
conn = sqlite3.connect(DB_PATH, timeout=30)
|
||||
conn.row_factory = sqlite3.Row
|
||||
conn.execute("PRAGMA journal_mode=WAL")
|
||||
conn.execute("PRAGMA busy_timeout=10000")
|
||||
return conn
|
||||
|
||||
|
||||
def run(workers: int = 8, limit: int = 200):
|
||||
"""Параллельная обработка очереди кандидатов. Возвращает (processed, classified, llm_ok, errors)."""
|
||||
stats = {"processed": 0, "classified": 0, "llm_ok": 0, "errors": 0, "lock": threading.Lock()}
|
||||
stop = threading.Event()
|
||||
|
||||
# каждый воркер сам создаёт соединение внутри своего потока
|
||||
print(f"workers={workers}, limit={limit}")
|
||||
with ThreadPoolExecutor(max_workers=workers) as ex:
|
||||
futs = {ex.submit(work_loop, i, stats, stop, limit): i
|
||||
for i in range(workers)}
|
||||
try:
|
||||
for f in as_completed(futs):
|
||||
f.result()
|
||||
except KeyboardInterrupt:
|
||||
stop.set()
|
||||
print("\nостановлено (Ctrl+C)")
|
||||
return stats["processed"], stats["classified"], stats["llm_ok"], stats["errors"]
|
||||
|
||||
|
||||
def main():
|
||||
ap = argparse.ArgumentParser(description="VESTI worker — параллельная обработка очереди кандидатов")
|
||||
ap.add_argument("--workers", type=int, default=8)
|
||||
ap.add_argument("--limit", type=int, default=200)
|
||||
args = ap.parse_args()
|
||||
|
||||
t0 = time.monotonic()
|
||||
proc, cls, llm, err = run(workers=args.workers, limit=args.limit)
|
||||
print(f"Обработано: {proc}, классифицировано: {cls}, через LLM: {llm}, ошибок: {err}, время: {time.monotonic()-t0:.1f}s", flush=True)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
Reference in New Issue
Block a user