# 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()