scraper: новая архитектура — core/ + forums/ (BaseAdapter, XenForo, phpBB)
- core/ — общий слой: fetch, throttle, output, state, runner - forums/base.py — абстрактный BaseAdapter - forums/xenforo/ — адаптер для XenForo (vwts.ru) - forums/phpbb/ — адаптер для phpBB (нива-лада.рф) - __main__.py — точка входа CLI - history/001-initial-structure.md — журнал изменений Также на ВМ (не в этом коммите): - lada_scraper/parse.py — пропуск sticky-тем при многопоточном парсинге
This commit is contained in:
@@ -0,0 +1,4 @@
|
||||
# scraper.core — общий слой парсера
|
||||
# Содержит модули, независимые от движка форума.
|
||||
# См. fetch.py, throttle.py, output.py, state.py, runner.py.
|
||||
|
||||
@@ -0,0 +1,118 @@
|
||||
"""HTTP-запросы: GET с ретраями.
|
||||
|
||||
Единая точка вызова HTTP для всех адаптеров.
|
||||
Не зависит от движка форума.
|
||||
"""
|
||||
|
||||
import time
|
||||
import logging
|
||||
|
||||
import requests
|
||||
from bs4 import BeautifulSoup
|
||||
|
||||
# ─── Модульный логгер ───────────────────────────────────────────────────
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
# ─── Глобальный throttle (устанавливается в runner.py) ──────────────────
|
||||
# fetch.get() вызывает throttle.wait() перед каждым запросом.
|
||||
throttle = None # type: ignore
|
||||
|
||||
|
||||
# ─── Дефолтные значения ─────────────────────────────────────────────────
|
||||
TIMEOUT = (10, 30) # (connect, read) — requests-таймаут
|
||||
HEADERS = {
|
||||
"User-Agent": (
|
||||
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
|
||||
"AppleWebKit/537.36 (KHTML, like Gecko) "
|
||||
"Chrome/125.0.0.0 Safari/537.36"
|
||||
),
|
||||
}
|
||||
RETRIES = 3 # максимум HTTP-попыток на один URL
|
||||
|
||||
|
||||
def make_session() -> requests.Session:
|
||||
"""Создать HTTP-сессию с дефолтными заголовками.
|
||||
|
||||
Returns:
|
||||
requests.Session с предустановленным User-Agent.
|
||||
"""
|
||||
session = requests.Session()
|
||||
session.headers.update(HEADERS)
|
||||
return session
|
||||
|
||||
|
||||
def get(url: str, session=None) -> BeautifulSoup | None:
|
||||
"""GET-запрос с ретраями, парсинг в BeautifulSoup (lxml).
|
||||
|
||||
Логика повторных попыток:
|
||||
1. 404 → сразу None (бесполезно ретраить).
|
||||
2. 403/429/503 → ждём 10*N сек, потом ретрай.
|
||||
3. Сеть/Таймаут → ждём 5 сек, потом ретрай.
|
||||
4. Прочее → ждём 5 сек, потом ретрай.
|
||||
5. Все RETRIES исчерпаны → None.
|
||||
|
||||
Args:
|
||||
url: Полный URL для загрузки.
|
||||
session: Сессия (если None — временная).
|
||||
|
||||
Returns:
|
||||
BeautifulSoup или None при недоступности.
|
||||
"""
|
||||
# ── Сессия ──────────────────────────────────────────────────────────
|
||||
if session is None:
|
||||
session = make_session()
|
||||
|
||||
# ── Throttle ────────────────────────────────────────────────────────
|
||||
if throttle is not None:
|
||||
throttle.wait()
|
||||
|
||||
# ── Цикл ретраев ────────────────────────────────────────────────────
|
||||
for attempt in range(1, RETRIES + 1):
|
||||
try:
|
||||
response = session.get(url, timeout=TIMEOUT)
|
||||
|
||||
# 404 — не найдено
|
||||
if response.status_code == 404:
|
||||
log.warning(" ⏭️ 404: %s", url[:80])
|
||||
return None
|
||||
|
||||
# 403/429/503 — сервер банит/перегружен
|
||||
if response.status_code in (403, 429, 503):
|
||||
wait = 10 * attempt
|
||||
log.warning(
|
||||
"⚠️ HTTP %d — попытка %d/%d, жду %d с",
|
||||
response.status_code, attempt, RETRIES, wait,
|
||||
)
|
||||
time.sleep(wait)
|
||||
continue
|
||||
|
||||
# Любая другая HTTP-ошибка
|
||||
response.raise_for_status()
|
||||
|
||||
# Успех
|
||||
return BeautifulSoup(response.text, "lxml")
|
||||
|
||||
except (requests.ConnectionError, requests.Timeout) as exc:
|
||||
log.warning(
|
||||
"⚠️ Сеть: %s — попытка %d/%d, жду 5 с",
|
||||
exc, attempt, RETRIES,
|
||||
)
|
||||
time.sleep(5)
|
||||
|
||||
except requests.HTTPError as exc:
|
||||
log.warning(
|
||||
"⚠️ HTTP %s — попытка %d/%d, жду 5 с",
|
||||
exc, attempt, RETRIES,
|
||||
)
|
||||
time.sleep(5)
|
||||
|
||||
except Exception as exc:
|
||||
log.warning(
|
||||
"⚠️ Ошибка: %s — попытка %d/%d, жду 5 с",
|
||||
exc, attempt, RETRIES,
|
||||
)
|
||||
time.sleep(5)
|
||||
|
||||
# Все попытки исчерпаны
|
||||
log.error("❌ Не загружено (%d попыток): %s", RETRIES, url[:80])
|
||||
return None
|
||||
@@ -0,0 +1,134 @@
|
||||
"""Запись данных: JSONL с ротацией по N тем на файл.
|
||||
|
||||
Формат одной записи (одна тема):
|
||||
{
|
||||
"section": str, # Название раздела
|
||||
"section_id": str, # ID раздела
|
||||
"topic_id": int, # ID темы на форуме
|
||||
"topic_title": str, # Заголовок темы
|
||||
"topic_url": str, # Прямая ссылка на тему
|
||||
"tags": [str], # Список тегов (если есть)
|
||||
"total_posts": int, # Количество постов в теме
|
||||
"posts": [ # Список постов
|
||||
{
|
||||
"author": str, # Имя автора
|
||||
"date": str, # Дата/время поста
|
||||
"num": str, # Номер поста (#1, #2...)
|
||||
"text": str, # Текст поста
|
||||
}
|
||||
],
|
||||
"scraped_at": str, # TIMESTAMP в ISO-формате
|
||||
}
|
||||
"""
|
||||
|
||||
import json
|
||||
import logging
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class TopicWriter:
|
||||
"""Потокобезопасная запись тем в JSONL-файлы с ротацией.
|
||||
|
||||
Каждый воркер пишет в свою папку (worker_N/).
|
||||
При достижении TOPICS_PER_FILE — создаётся новый part_*.jsonl.
|
||||
|
||||
Атрибуты:
|
||||
_worker_dir: Путь к папке воркера.
|
||||
_topics_file: Счётчик записанных тем в текущий файл.
|
||||
_file_index: Индекс текущего part_*.jsonl.
|
||||
_handle: Открытый файловый дескриптор.
|
||||
"""
|
||||
|
||||
def __init__(self, worker_dir: Path, topics_per_file: int = 100):
|
||||
"""Инициализация.
|
||||
|
||||
Args:
|
||||
worker_dir: Папка воркера (worker_N/).
|
||||
topics_per_file: Сколько тем писать в один файл перед ротацией.
|
||||
"""
|
||||
self._worker_dir = worker_dir
|
||||
self._topics_per_file = topics_per_file
|
||||
|
||||
# Счётчики
|
||||
self._topics_file = 0
|
||||
self._file_index = 0
|
||||
|
||||
# Открываем первый файл
|
||||
self._handle = self._open_next()
|
||||
|
||||
# ── Открытие/закрытие файла ─────────────────────────────────────────
|
||||
|
||||
def _open_next(self):
|
||||
"""Закрыть текущий файл (если есть) и открыть следующий part_N.jsonl."""
|
||||
# Закрываем предыдущий
|
||||
if hasattr(self, "_handle") and self._handle and not self._handle.closed:
|
||||
self._handle.close()
|
||||
|
||||
# Новый путь
|
||||
path = self._worker_dir / f"part_{self._file_index:04d}.jsonl"
|
||||
|
||||
# Сбрасываем счётчик тем в файле
|
||||
self._topics_file = 0
|
||||
|
||||
log.info(" 💾 %s", path.name)
|
||||
fh = open(path, "w", encoding="utf-8")
|
||||
self._file_index += 1
|
||||
return fh
|
||||
|
||||
def close(self):
|
||||
"""Закрыть текущий файл."""
|
||||
if self._handle and not self._handle.closed:
|
||||
self._handle.close()
|
||||
|
||||
# ── Запись ──────────────────────────────────────────────────────────
|
||||
|
||||
def write(self,
|
||||
section_name: str,
|
||||
section_slug: str,
|
||||
topic: dict,
|
||||
posts: list,
|
||||
tags: list | None = None):
|
||||
"""Записать одну тему в JSONL.
|
||||
|
||||
Если текущий файл заполнен — автоматически создаёт новый.
|
||||
|
||||
Args:
|
||||
section_name: Название раздела (напр. "Бензиновые двигатели").
|
||||
section_slug: Короткий ID раздела (напр. "benzinovye-dvigateli").
|
||||
topic: Словарь темы: {id, title, url}.
|
||||
posts: Список постов: [{author, date, num, text}].
|
||||
tags: Список тегов (опционально).
|
||||
"""
|
||||
# ── Ротация ─────────────────────────────────────────────────────
|
||||
if self._topics_file >= self._topics_per_file:
|
||||
self._open_next()
|
||||
|
||||
# ── Формируем запись ────────────────────────────────────────────
|
||||
record = {
|
||||
"section": section_name,
|
||||
"section_id": section_slug,
|
||||
"topic_id": topic["id"],
|
||||
"topic_title": topic["title"],
|
||||
"topic_url": topic.get("url", ""),
|
||||
"tags": tags or [],
|
||||
"total_posts": len(posts),
|
||||
"posts": posts,
|
||||
"scraped_at": datetime.now(timezone.utc).isoformat(),
|
||||
}
|
||||
|
||||
# ── Пишем в файл ────────────────────────────────────────────────
|
||||
line = json.dumps(record, ensure_ascii=False)
|
||||
self._handle.write(line + "\n")
|
||||
self._handle.flush()
|
||||
self._topics_file += 1
|
||||
|
||||
# ── Контекстный менеджер ────────────────────────────────────────────
|
||||
|
||||
def __enter__(self):
|
||||
return self
|
||||
|
||||
def __exit__(self, *args):
|
||||
self.close()
|
||||
@@ -0,0 +1,322 @@
|
||||
"""Главный цикл: обход разделов через адаптер.
|
||||
|
||||
Единая точка запуска для любого движка форума.
|
||||
Не содержит кода парсинга — использует адаптер из forums/.
|
||||
"""
|
||||
|
||||
import sys
|
||||
import logging
|
||||
import shutil
|
||||
import json
|
||||
from concurrent.futures import ThreadPoolExecutor, as_completed
|
||||
from pathlib import Path
|
||||
|
||||
from . import fetch as _fetch
|
||||
from .throttle import Throttle, BanDetector
|
||||
from .output import TopicWriter
|
||||
from .state import load, save, merge_workers
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
# ─── Глобальный throttle (виден fetch.py) ───────────────────────────────
|
||||
_throttle: Throttle | None = None
|
||||
|
||||
|
||||
# ═══════════════════════════════════════════════════════════════════════════
|
||||
# Функция одного воркера (запускается в ThreadPoolExecutor)
|
||||
# ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
def _worker(section_url: str,
|
||||
section_slug: str,
|
||||
section_name: str,
|
||||
pages: range,
|
||||
worker_id: int,
|
||||
adapter) -> tuple:
|
||||
"""Обойти свой диапазон страниц и сохранить темы в worker_N/.
|
||||
|
||||
Args:
|
||||
section_url: URL раздела (первая страница).
|
||||
section_slug: Короткое имя для папки (напр. "diagnostika").
|
||||
section_name: Человеческое название раздела.
|
||||
pages: Диапазон страниц для этого воркера.
|
||||
worker_id: Номер воркера (0, 1, 2...).
|
||||
adapter: Экземпляр адаптера (forums/*/adapter.py).
|
||||
|
||||
Returns:
|
||||
(worker_id, topic_count) — для汇总 в главном потоке.
|
||||
"""
|
||||
session = _fetch.make_session()
|
||||
|
||||
# ── Состояние воркера ───────────────────────────────────────────────
|
||||
worker_state = {
|
||||
"topics_done": {},
|
||||
"pages_done": [],
|
||||
"file_index": 0,
|
||||
"topic_count": 0,
|
||||
}
|
||||
|
||||
consecutive_404 = 0
|
||||
writer = None
|
||||
|
||||
# ── Обход страниц раздела ───────────────────────────────────────────
|
||||
for page_num in pages:
|
||||
# Формируем URL страницы через адаптер
|
||||
url = adapter.page_url(section_url, page_num)
|
||||
log.info("[W%d] 📄 стр.%d: %s", worker_id, page_num, url[:100])
|
||||
|
||||
# Загружаем
|
||||
soup = _fetch.get(url, session)
|
||||
if soup is None:
|
||||
# 404 подряд — вероятно, страницы кончились
|
||||
consecutive_404 += 1
|
||||
if consecutive_404 >= 3:
|
||||
log.info("[W%d] ⏹️ 3 пустых страницы подряд — завершаем", worker_id)
|
||||
break
|
||||
continue
|
||||
consecutive_404 = 0
|
||||
|
||||
# Парсим список тем через адаптер
|
||||
topics = adapter.parse_topics(soup)
|
||||
log.info("[W%d] Тем на странице: %d", worker_id, len(topics))
|
||||
|
||||
# ── Обход тем на странице ────────────────────────────────────────
|
||||
for idx, topic in enumerate(topics, 1):
|
||||
topic_id = str(topic["id"])
|
||||
|
||||
# Пропускаем уже обработанные
|
||||
if topic_id in worker_state["topics_done"]:
|
||||
continue
|
||||
|
||||
log.info(
|
||||
"[W%d] [%d/%d] #%s: %s",
|
||||
worker_id, idx, len(topics), topic_id, topic["title"][:80],
|
||||
)
|
||||
|
||||
# Загружаем страницу темы (первую)
|
||||
topic_url = adapter.topic_url(section_url, topic)
|
||||
topic["url"] = topic_url
|
||||
topic_soup = _fetch.get(topic_url, session)
|
||||
if topic_soup is None:
|
||||
# Не загрузилась — помечаем как сделанную и идём дальше
|
||||
worker_state["topics_done"][topic_id] = topic["title"]
|
||||
continue
|
||||
|
||||
# Определяем количество страниц в теме
|
||||
topic_pages = adapter.count_pages(topic_soup)
|
||||
log.info(
|
||||
"[W%d] Страниц в теме: %d",
|
||||
worker_id, topic_pages,
|
||||
)
|
||||
|
||||
# Парсим посты с первой страницы
|
||||
all_posts = adapter.parse_posts(topic_soup)
|
||||
|
||||
# Догружаем остальные страницы темы
|
||||
for tp in range(2, topic_pages + 1):
|
||||
tp_url = adapter.topic_page_url(topic_url, tp)
|
||||
tp_soup = _fetch.get(tp_url, session)
|
||||
if tp_soup:
|
||||
all_posts.extend(adapter.parse_posts(tp_soup))
|
||||
|
||||
# Парсим теги (если адаптер поддерживает)
|
||||
tags = adapter.parse_tags(topic_soup) if hasattr(adapter, "parse_tags") else []
|
||||
|
||||
log.info(
|
||||
"[W%d] %d постов, теги: %s",
|
||||
worker_id, len(all_posts), tags,
|
||||
)
|
||||
|
||||
# ── Сохраняем ────────────────────────────────────────────────
|
||||
# Лениво инициализируем writer при первой теме
|
||||
if writer is None:
|
||||
worker_dir = Path(adapter.data_dir) / section_slug / f"worker_{worker_id}"
|
||||
worker_dir.mkdir(parents=True, exist_ok=True)
|
||||
writer = TopicWriter(worker_dir, adapter.topics_per_file)
|
||||
|
||||
writer.write(section_name, section_slug, topic, all_posts, tags)
|
||||
worker_state["topics_done"][topic_id] = topic["title"]
|
||||
worker_state["topic_count"] += 1
|
||||
|
||||
# Помечаем страницу как пройденную
|
||||
worker_state["pages_done"].append(page_num)
|
||||
|
||||
# ── Сохраняем состояние воркера ─────────────────────────────────────
|
||||
if writer is not None:
|
||||
writer.close()
|
||||
|
||||
worker_dir = Path(adapter.data_dir) / section_slug / f"worker_{worker_id}"
|
||||
worker_dir.mkdir(parents=True, exist_ok=True)
|
||||
state_file = worker_dir / "state.json"
|
||||
state_file.write_text(
|
||||
json.dumps(worker_state, ensure_ascii=False, indent=2),
|
||||
encoding="utf-8",
|
||||
)
|
||||
|
||||
session.close()
|
||||
return worker_id, worker_state["topic_count"]
|
||||
|
||||
|
||||
# ═══════════════════════════════════════════════════════════════════════════
|
||||
# Функция мержа воркеров
|
||||
# ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
def _merge_workers(adapter, section_slug: str):
|
||||
"""Собрать part_*.jsonl из worker_N/ в корень раздела.
|
||||
|
||||
Переименовывает файлы с единой нумерацией, удаляет пустые
|
||||
worker-папки, мержит state.json.
|
||||
|
||||
Args:
|
||||
adapter: Экземпляр адаптера.
|
||||
section_slug: Короткое имя раздела.
|
||||
"""
|
||||
root = Path(adapter.data_dir) / section_slug
|
||||
|
||||
# ── Собираем все part_*.jsonl из worker_N/ ──────────────────────────
|
||||
all_parts = []
|
||||
for worker_dir in sorted(root.glob("worker_*")):
|
||||
if worker_dir.is_dir():
|
||||
for part_file in sorted(worker_dir.glob("part_*.jsonl")):
|
||||
all_parts.append(part_file)
|
||||
|
||||
if not all_parts:
|
||||
log.warning("⚠️ Нет файлов для мержа в %s", root)
|
||||
return
|
||||
|
||||
log.info(
|
||||
"🔗 Мерж %d файлов из %d воркеров → %s",
|
||||
len(all_parts), len(list(root.glob("worker_*"))), root,
|
||||
)
|
||||
|
||||
# ── Переименовываем с единой нумерацией ─────────────────────────────
|
||||
for new_index, src in enumerate(all_parts):
|
||||
dst = root / f"part_{new_index:04d}.jsonl"
|
||||
shutil.move(str(src), str(dst))
|
||||
|
||||
# ── Удаляем пустые папки воркеров ───────────────────────────────────
|
||||
for worker_dir in root.glob("worker_*"):
|
||||
if worker_dir.is_dir():
|
||||
try:
|
||||
worker_dir.rmdir() # удаляет только пустую папку
|
||||
except OSError:
|
||||
pass # если не пустая — оставляем
|
||||
|
||||
# ── Мержим state.json ───────────────────────────────────────────────
|
||||
merge_workers(adapter.data_dir, section_slug)
|
||||
|
||||
log.info("✅ Мерж завершён: %d файлов → %s", len(all_parts), root)
|
||||
|
||||
|
||||
# ═══════════════════════════════════════════════════════════════════════════
|
||||
# Публичный запуск
|
||||
# ═══════════════════════════════════════════════════════════════════════════
|
||||
|
||||
def run(adapter,
|
||||
section_url: str,
|
||||
section_slug: str | None = None,
|
||||
workers: int = 4,
|
||||
no_delay: bool = False,
|
||||
ban_test: bool = True):
|
||||
"""Главный цикл: многопоточный обход раздела через адаптер.
|
||||
|
||||
Порядок работы:
|
||||
1. Проверка бана (BanDetector).
|
||||
2. Загрузка первой страницы для определения названия и пагинации.
|
||||
3. Разбиение страниц на диапазоны по числу воркеров.
|
||||
4. Запуск ThreadPoolExecutor.
|
||||
5. Мерж результатов.
|
||||
|
||||
Args:
|
||||
adapter: Экземпляр адаптера (forums/*/adapter.py).
|
||||
section_url: URL раздела (первая страница).
|
||||
section_slug: Короткое имя (если None — берётся из URL).
|
||||
workers: Количество потоков.
|
||||
no_delay: True — отключить координацию задержек.
|
||||
ban_test: True — проверить бан перед стартом.
|
||||
"""
|
||||
# ── Определяем slug из URL ──────────────────────────────────────────
|
||||
if section_slug is None:
|
||||
section_slug = section_url.rstrip("/").split("/")[-1]
|
||||
|
||||
# ── Настройка throttle ──────────────────────────────────────────────
|
||||
throttle_enabled = not no_delay
|
||||
|
||||
# Если workers=1, координация не нужна — один поток не конфликтует
|
||||
if workers <= 1:
|
||||
throttle_enabled = False
|
||||
|
||||
# Проверка бана
|
||||
if ban_test and throttle_enabled and workers > 1:
|
||||
if BanDetector.test(section_url, adapter.ban_test_count):
|
||||
throttle_enabled = True
|
||||
else:
|
||||
throttle_enabled = False
|
||||
|
||||
# Устанавливаем глобальный throttle
|
||||
global _throttle
|
||||
_throttle = Throttle(adapter.delay, enabled=throttle_enabled)
|
||||
_fetch.throttle = _throttle
|
||||
|
||||
log.info(
|
||||
"⚙️ Потоков: %d, координация: %s",
|
||||
workers, "вкл" if throttle_enabled else "выкл",
|
||||
)
|
||||
|
||||
# ── Определяем название раздела и число страниц ─────────────────────
|
||||
session = _fetch.make_session()
|
||||
soup = _fetch.get(section_url, session)
|
||||
if soup is None:
|
||||
log.error("❌ Сайт недоступен: %s", section_url)
|
||||
sys.exit(1)
|
||||
|
||||
# Название раздела (из h1 или из URL)
|
||||
section_name = adapter.parse_section_name(soup) or section_slug
|
||||
|
||||
# Количество страниц в разделе
|
||||
total_pages = adapter.count_pages(soup)
|
||||
log.info("📋 %s | страниц: %d", section_name, total_pages)
|
||||
session.close()
|
||||
|
||||
# ── Разбиваем страницы на диапазоны для воркеров ────────────────────
|
||||
if workers > total_pages:
|
||||
workers = max(1, total_pages)
|
||||
|
||||
base_size = total_pages // workers
|
||||
remainder = total_pages % workers
|
||||
|
||||
ranges = []
|
||||
start = 1
|
||||
for w in range(workers):
|
||||
size = base_size + (1 if w < remainder else 0)
|
||||
ranges.append(range(start, start + size))
|
||||
start += size
|
||||
|
||||
log.info("📦 Диапазоны: %s", [f"{r.start}-{r[-1]}" for r in ranges])
|
||||
|
||||
# ── Запуск воркеров ─────────────────────────────────────────────────
|
||||
total_topics = 0
|
||||
with ThreadPoolExecutor(max_workers=workers) as pool:
|
||||
futures = [
|
||||
pool.submit(_worker, section_url, section_slug, section_name,
|
||||
rng, wid, adapter)
|
||||
for wid, rng in enumerate(ranges)
|
||||
]
|
||||
for future in as_completed(futures):
|
||||
wid, cnt = future.result()
|
||||
total_topics += cnt
|
||||
log.info("[W%d] ✅ завершён: %d тем", wid, cnt)
|
||||
|
||||
# ── Мерж ────────────────────────────────────────────────────────────
|
||||
_merge_workers(adapter, section_slug)
|
||||
|
||||
# ── Итог ────────────────────────────────────────────────────────────
|
||||
data_dir = Path(adapter.data_dir)
|
||||
part_files = sorted((data_dir / section_slug).glob("part_*.jsonl"))
|
||||
total_size_mb = sum(f.stat().st_size for f in part_files) / (1024 * 1024)
|
||||
|
||||
log.info("=" * 60)
|
||||
log.info(
|
||||
"✅ '%s' — %d тем в %d частях (%.1f MB)",
|
||||
section_name, total_topics, len(part_files), total_size_mb,
|
||||
)
|
||||
log.info("=" * 60)
|
||||
@@ -0,0 +1,114 @@
|
||||
"""Управление состоянием: state.json в папке раздела.
|
||||
|
||||
Позволяет:
|
||||
- Сохранять прогресс (какие темы обработаны, какие страницы пройдены).
|
||||
- Продолжать парсинг с места остановки (той же командой).
|
||||
- Мержить состояния воркеров в единый state.json (после завершения).
|
||||
|
||||
Файл state.json (раздел):
|
||||
{
|
||||
"topics_done": {"123": "Название темы", "456": "..."},
|
||||
"pages_done": [1, 2, 3],
|
||||
"file_index": 5,
|
||||
"topic_count": 432
|
||||
}
|
||||
|
||||
Файл worker_N/state.json:
|
||||
Аналогичная структура, но только для одного воркера.
|
||||
"""
|
||||
|
||||
import json
|
||||
import logging
|
||||
from pathlib import Path
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def _state_file(data_dir: Path, section_slug: str) -> Path:
|
||||
"""Полный путь к state.json для раздела.
|
||||
|
||||
Args:
|
||||
data_dir: Корневая папка с данными (напр. ./vwts_data).
|
||||
section_slug: Короткое имя раздела.
|
||||
|
||||
Returns:
|
||||
Path к state.json.
|
||||
"""
|
||||
section_dir = data_dir / section_slug
|
||||
section_dir.mkdir(parents=True, exist_ok=True)
|
||||
return section_dir / "state.json"
|
||||
|
||||
|
||||
def load(data_dir: Path, section_slug: str) -> dict:
|
||||
"""Загрузить state.json раздела.
|
||||
|
||||
Args:
|
||||
data_dir: Корневая папка с данными.
|
||||
section_slug: Короткое имя раздела.
|
||||
|
||||
Returns:
|
||||
Словарь состояния. Если файла нет — пустой state.
|
||||
"""
|
||||
path = _state_file(data_dir, section_slug)
|
||||
|
||||
if path.exists():
|
||||
try:
|
||||
return json.loads(path.read_text(encoding="utf-8"))
|
||||
except (json.JSONDecodeError, OSError) as exc:
|
||||
log.warning("⚠️ Ошибка чтения %s: %s", path, exc)
|
||||
|
||||
# Пустой state по умолчанию
|
||||
return {"topics_done": {}, "pages_done": [], "file_index": 0, "topic_count": 0}
|
||||
|
||||
|
||||
def save(data_dir: Path, section_slug: str, state: dict):
|
||||
"""Сохранить state.json раздела.
|
||||
|
||||
Args:
|
||||
data_dir: Корневая папка с данными.
|
||||
section_slug: Короткое имя раздела.
|
||||
state: Словарь состояния для сохранения.
|
||||
"""
|
||||
path = _state_file(data_dir, section_slug)
|
||||
path.write_text(
|
||||
json.dumps(state, ensure_ascii=False, indent=2),
|
||||
encoding="utf-8",
|
||||
)
|
||||
|
||||
|
||||
def merge_workers(data_dir: Path, section_slug: str):
|
||||
"""Собрать состояния всех воркеров в единый state.json раздела.
|
||||
|
||||
Проходит по worker_*/state.json, сливает topics_done, pages_done,
|
||||
суммирует topic_count.
|
||||
|
||||
Args:
|
||||
data_dir: Корневая папка с данными.
|
||||
section_slug: Короткое имя раздела.
|
||||
"""
|
||||
root = data_dir / section_slug
|
||||
|
||||
# ── Собираем все worker_*/state.json ────────────────────────────────
|
||||
merged = {"topics_done": {}, "pages_done": [], "file_index": 0, "topic_count": 0}
|
||||
|
||||
for worker_state_path in sorted(root.glob("worker_*/state.json")):
|
||||
try:
|
||||
worker_state = json.loads(
|
||||
worker_state_path.read_text(encoding="utf-8"),
|
||||
)
|
||||
merged["topics_done"].update(worker_state.get("topics_done", {}))
|
||||
merged["pages_done"].extend(worker_state.get("pages_done", []))
|
||||
merged["topic_count"] += worker_state.get("topic_count", 0)
|
||||
except (json.JSONDecodeError, OSError) as exc:
|
||||
log.warning("⚠️ Ошибка чтения %s: %s", worker_state_path, exc)
|
||||
|
||||
# Чистим дубликаты и сортируем
|
||||
merged["pages_done"] = sorted(set(merged["pages_done"]))
|
||||
|
||||
# Определяем следующий file_index по количеству part_*.jsonl
|
||||
existing_parts = list(root.glob("part_*.jsonl"))
|
||||
merged["file_index"] = max(0, len(existing_parts) - 1)
|
||||
|
||||
# ── Пишем ───────────────────────────────────────────────────────────
|
||||
save(data_dir, section_slug, merged)
|
||||
log.info("🔗 Мерж состояний: %d воркеров → state.json", len(list(root.glob("worker_*"))))
|
||||
@@ -0,0 +1,147 @@
|
||||
"""Координация задержек между потоками + автоопределение бана.
|
||||
|
||||
Throttle:
|
||||
Thread-safe координатор: максимум 1 запрос в DELAY секунд.
|
||||
При отключённой координации — просто time.sleep(DELAY).
|
||||
|
||||
BanDetector:
|
||||
Проверяет, банит ли сервер за быстрые запросы.
|
||||
Шлёт COUNT запросов подряд без задержки, ищет 403/429.
|
||||
"""
|
||||
|
||||
import time
|
||||
import threading
|
||||
import logging
|
||||
|
||||
import requests
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class Throttle:
|
||||
"""Координатор задержек между потоками.
|
||||
|
||||
Позволяет запускать N потоков, но реальные HTTP-запросы
|
||||
идут не чаще 1 в delay секунд. Потоки ждут в очереди
|
||||
через threading.Lock.
|
||||
|
||||
Атрибуты:
|
||||
_delay: Минимальный интервал между запросами (сек).
|
||||
_enabled: True — координация включена, False — только sleep.
|
||||
_lock: Мьютекс для потокобезопасности.
|
||||
_last: Монотонное время последнего запроса.
|
||||
"""
|
||||
|
||||
def __init__(self, delay: float = 1.5, enabled: bool = True):
|
||||
"""Инициализация.
|
||||
|
||||
Args:
|
||||
delay: Минимальный интервал между запросами (сек).
|
||||
enabled: Если False — вместо координации просто sleep(delay).
|
||||
"""
|
||||
self._delay = delay
|
||||
self._enabled = enabled
|
||||
self._lock = threading.Lock()
|
||||
self._last = 0.0 # time.monotonic() последнего запроса
|
||||
|
||||
# ── Свойства ────────────────────────────────────────────────────────
|
||||
|
||||
@property
|
||||
def delay(self) -> float:
|
||||
"""Текущая задержка."""
|
||||
return self._delay
|
||||
|
||||
@property
|
||||
def enabled(self) -> bool:
|
||||
"""Включена ли координация."""
|
||||
return self._enabled
|
||||
|
||||
# ── Основной метод ──────────────────────────────────────────────────
|
||||
|
||||
def wait(self):
|
||||
"""Подождать своей очереди на запрос.
|
||||
|
||||
Если координация ВЫКЛЮЧЕНА:
|
||||
Просто спим delay секунд.
|
||||
Если координация ВКЛЮЧЕНА:
|
||||
Захватываем мьютекс, считаем сколько прошло с _last,
|
||||
если меньше delay — спим разницу, обновляем _last.
|
||||
"""
|
||||
if not self._enabled:
|
||||
# Без координации: каждый поток спит свою задержку
|
||||
time.sleep(self._delay)
|
||||
return
|
||||
|
||||
with self._lock:
|
||||
elapsed = time.monotonic() - self._last
|
||||
if elapsed < self._delay:
|
||||
time.sleep(self._delay - elapsed)
|
||||
self._last = time.monotonic()
|
||||
|
||||
# ── Управление ──────────────────────────────────────────────────────
|
||||
|
||||
def disable(self):
|
||||
"""Отключить координацию."""
|
||||
self._enabled = False
|
||||
|
||||
def enable(self):
|
||||
"""Включить координацию."""
|
||||
self._enabled = True
|
||||
|
||||
|
||||
class BanDetector:
|
||||
"""Определяет, банит ли сервер за быстрые запросы.
|
||||
|
||||
Статический метод test() шлёт count запросов подряд
|
||||
(без задержек) и смотрит ответы. Если хоть один 403/429 —
|
||||
сервер банит.
|
||||
"""
|
||||
|
||||
@staticmethod
|
||||
def test(base_url: str, count: int = 5) -> bool:
|
||||
"""Проверить бан быстрыми запросами.
|
||||
|
||||
Args:
|
||||
base_url: Любой URL того же хоста (главная раздела).
|
||||
count: Сколько запросов сделать.
|
||||
|
||||
Returns:
|
||||
True — сервер банит (есть 403/429),
|
||||
False — бан не обнаружен.
|
||||
"""
|
||||
log.info("🔍 Проверка бана: %d запросов подряд...", count)
|
||||
|
||||
session = requests.Session()
|
||||
session.headers.update({
|
||||
"User-Agent": (
|
||||
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
|
||||
"AppleWebKit/537.36"
|
||||
),
|
||||
})
|
||||
|
||||
banned = False
|
||||
|
||||
for i in range(1, count + 1):
|
||||
try:
|
||||
response = session.get(base_url, timeout=(10, 30))
|
||||
log.info(" [%d/%d] HTTP %d", i, count, response.status_code)
|
||||
|
||||
# Сервер явно банит
|
||||
if response.status_code in (403, 429):
|
||||
banned = True
|
||||
break
|
||||
|
||||
except Exception as exc:
|
||||
# Сетевая ошибка — вероятно, бан на уровне соединения
|
||||
log.warning(" [%d/%d] ошибка: %s", i, count, exc)
|
||||
banned = True
|
||||
break
|
||||
|
||||
session.close()
|
||||
|
||||
if banned:
|
||||
log.warning("⚠️ Сервер банит быстрые запросы — включаю координацию")
|
||||
else:
|
||||
log.info("✅ Бан не обнаружен — без координации")
|
||||
|
||||
return banned
|
||||
Reference in New Issue
Block a user