mirror of
https://gitverse.ru/kpa39l/vesti.git
synced 2026-09-29 09:55:03 +00:00
170 lines
8.1 KiB
Python
170 lines
8.1 KiB
Python
# 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() |