gendesign/tradein-mvp/scripts/local-avito-msk/collect.py
bot-backend 8fcef103c2 feat(msk-collector): Циан как вторая площадка сбора по Москве и МО
Сборщик получает ключ --platform {avito,cian}: платформо-зависимые куски
(URL коридора, счётчик, парс, потолок пагинации, целевая таблица) вынесены
в PlatformAdapter, общая часть — бисекция по цене, guard на блок, накопитель
и заливка — одна на обе площадки.

Скоуп Циан проверен живыми запросами 10.09, а не взят из документации:

- region=1 и region=4593 двумя параметрами НЕ объединяются, сервер берёт
  последний. Прежний базовый URL собрал бы одну Московскую область и молча
  потерял Москву целиком (92 817 объявлений). Правильный скоуп — region=-1.
- object_type[0]=1 обязателен: без него счётчик считает новостройки, которые
  дальше не сохраняются, и расходится с числом карточек вдвое.
- Итоговый скоуп: 62 548 объявлений вторички по Москве и МО.
- Потолок пагинации 54 страницы подтверждён: страница 55 пуста.

Фильтра новостроек в адаптере нет, и это сознательное отличие от кита. Кит
считает новостройкой всё, у чего есть offer.newbuilding.id (serp.py:401-402),
потому что для ЕКБ выдача бралась без object_type. На московской странице из
28 карточек 16 имеют этот блок, и у всех шестнадцати isFromBuilder=false,
isFromLeadFactory=false — это вторичка в ЖК, а не лоты застройщика. С китовым
фильтром прогон терял бы 57 % корпуса безвозвратно. msk_raw — сырьё, payload
несёт listing_segment целиком, сегмент отделяется на импорте в listings.

Детект блока Циан. Капча приходит как HTTP 200 с обычной на вид страницей:
поймана живьём, 40 КБ, title «Captcha - база объявлений ЦИАН», без
window._cianConfig. По статусу её не отличить, поэтому маркер ищется в первых
4 КБ (в нормальной выдаче на 2,6 МБ там ни одного вхождения). Вторая сеть в
parse_page: счётчик None при нуле сырых карточек — тоже стоп. Без этого блок
засчитывался бы как пустая страница, коридор уходил в done, а --resume его
уже не перебрал бы.

Гвард empty_page считает карточки ДО фильтра: иначе штатный ноль после
фильтрации был бы неотличим от блока. У Авито атрибута нет, дефолт — длина
итогового списка, поведение прежнее.

Заливка. Целевая таблица берётся из адаптера, staging создаётся с
INCLUDING IDENTITY (без него identity-колонка не переносится, а NOT NULL
переносится всегда, и \copy падал бы на каждом батче). Наличие таблицы
проверяется на старте одним SELECT to_regclass — иначе прогон умирал бы
на первой заливке, после часов планирования.

Миграция 299 заводит cian_cards, domclick_cards, yandex_cards и три вью
по образцу avito_cards/avito_latest. id — bigserial, а не IDENTITY, ровно
по причине выше.

Заодно: _wt-mskcol/ выведен из индекса и закрыт в .gitignore. Копия-worktree
попала в репозиторий 09.09, и три фикса ушли в дубликат мимо канонического
collect.py — дефект нашёлся только при сверке размера страницы.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01VQ8jqr4SFirX5tFLwdSrXh
2026-09-10 13:23:46 +03:00

1066 lines
52 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#!/usr/bin/env python3
"""Локальный ручной сборщик SERP Авито/Циан по Москве и МО (эпик #2989, трек 1).
Запускается ВРУЧНУЮ с машины владельца. Прод-скрейпер, его расписания и
прокси-пул не задействованы вообще: браузер — уже открытый Chrome владельца
(подключение по CDP), парсер — импорт из scraper-kit, заливка — поток в psql
через ssh. Скрипт ничего не устанавливает и своего профиля не поднимает.
Платформа выбирается ключом --platform {avito,cian} (дефолт avito) — см. класс
PlatformAdapter ниже. У каждой платформы свой потолок пагинации, свой билдер
URL коридора и своя целевая таблица в msk_raw.
Дефолтный режим — --measure 100 (замер): полный проход только по явному --full.
"""
from __future__ import annotations
import argparse
import asyncio
import csv
import io
import json
import math
import os
import random
import re
import shlex
import subprocess
import sys
import time
from dataclasses import dataclass, field
from datetime import datetime, timezone
from pathlib import Path
from types import SimpleNamespace
from typing import Any, Callable, Iterable
from urllib.parse import parse_qsl, urlencode, urlsplit, urlunsplit
# --- импорт парсеров из scraper-kit без установки backend -------------------
_KIT_SRC = Path(__file__).resolve().parents[2] / "packages" / "scraper-kit" / "src"
if str(_KIT_SRC) not in sys.path:
sys.path.insert(0, str(_KIT_SRC))
from scraper_kit.providers.avito.serp import ( # noqa: E402
AvitoScraper,
_is_firewall_page,
)
from scraper_kit.providers.cian.serp import ( # noqa: E402
CianScraper,
_CIAN_OFFERS_PER_PAGE,
)
# Вкладка, открытая у владельца: вторичка, Москва + МО.
DEFAULT_AVITO_BASE_URL = (
"https://www.avito.ru/moskva_i_mo/kvartiry/prodam/vtorichka-ASgBAgICAkSSA8YQ5geMUg"
"?f=ASgBAgICA0SSA8YQ5geMUuDC3c0CgOkP"
)
# domain www.cian.ru, а не ekb.cian.ru: у CianScraper._build_url домен захардкожен
# под ЕКБ (self.base_url класс-константа), сюда не подходит — поэтому URL
# коридора строим своим билдером (_cian_build_url), а не scraper._build_url.
#
# Параметры проверены живым запросом 10.09 (curl, <title> и totalOffers из SSR):
# region=1 -> «Купить квартиру в Москве — 92 817»
# region=4593 -> «Купить квартиру в Московской области — 60 231»
# region=1&region=4593 -> «Московская область — 60 231»: побеждает ПОСЛЕДНИЙ,
# объединения по нескольким region= НЕТ
# region=-1 -> «Москва и Московская область — 153 049»
# region=-1 + object_type[0]=1 -> «Москва и МО — 62 548», puid3=sale_type_second
# region=1 + object_type[0]=1 -> «Москва — 36 743»
# Отсюда region=-1 — единственный способ получить единый скоуп Москва+МО.
# object_type[0]=1 (вторичка) обязателен: без него totalOffers считает и
# новостройки, которые парсер выбрасывает, — счётчик и карточки расходятся вдвое.
DEFAULT_CIAN_BASE_URL = (
"https://www.cian.ru/cat.php?deal_type=sale&engine_version=2&offer_type=flat"
"&region=-1&object_type%5B0%5D=1&sort=creation_date_desc"
)
# Замерено живым проходом (не из документации): выдача Москва+МО отдаёт 50 карточек
# на страницу. Пока здесь стояло 60, planned_pages считал count/60 и не запрашивал
# последние ~17% каждого коридора — 34 753 по счётчику против 28 352 собранных.
# Молчаливое усечение читается как «покрыто всё», поэтому число проверяется живьём.
AVITO_PAGE_SIZE = 50 # карточек на странице выдачи Авито
AVITO_MAX_PAGES = 30 # потолок пагинации Авито → 30*50 = 1500 на один запрос
# Циан отдаёт ~28 офферов/страницу (_CIAN_OFFERS_PER_PAGE в providers/cian/serp.py).
# Потолок страниц на бакет — 54, см. CianScraper._paginate_leaf_bucket
# (max_pages_per_bucket: int = 54, комментарий "Cian hard cap ~54; не превышать").
# 28*54 = 1512 — заметно меньше авитовских 1800, поэтому свой HARD_CAP обязателен.
CIAN_PAGE_SIZE = _CIAN_OFFERS_PER_PAGE
CIAN_MAX_PAGES = 54
PRICE_FLOOR = 500_000 # нижняя граница осмысленного коридора, ₽
PRICE_PROBE_START = 8_000_000 # старт удвоения при поиске верхней границы
PRICE_CEIL = 2_000_000_000
MIN_WIDTH_RATIO = 1.05 # уже этого коридор не делим (геометрическая ширина)
MAX_DEPTH = 12
_BATCH_ID_RE = re.compile(r"^[A-Za-z0-9._-]+$")
_POW_MARKERS = ("startpow", "доступ ограничен: проверка безопасности")
# Заполняется при входе в Loader.__aenter__ (playwright импортируется лениво,
# чтобы --help работал без установленного пакета).
PlaywrightTimeoutError: type[BaseException] = TimeoutError
# Через столько загрузок вкладка сборщика пересоздаётся (см. _recycle_if_needed).
PAGE_RECYCLE_EVERY = 75
class Blocked(Exception):
"""Первый признак блока. Ретраев нет — только немедленный стоп."""
def __init__(self, reason: str, detail: str = "") -> None:
super().__init__(f"{reason}: {detail}" if detail else reason)
self.reason = reason
self.detail = detail
class BudgetExhausted(Exception):
"""Потолок --measure выбран: штатный выход, не ошибка."""
# --- план коридоров --------------------------------------------------------
@dataclass
class Corridor:
lo: int | None
hi: int | None
count: int | None = None
truncated: bool = False
pages_done: int = 0
status: str = "pending" # pending | done
missed: int = 0 # заведомо недобрано (count - HARD_CAP), если truncated
def label(self) -> str:
lo = "-" if self.lo is None else f"{self.lo:_}"
hi = "-" if self.hi is None else f"{self.hi:_}"
return f"[{lo} .. {hi}]"
def to_json(self) -> dict[str, Any]:
return {
"lo": self.lo, "hi": self.hi, "count": self.count,
"truncated": self.truncated, "pages_done": self.pages_done,
"status": self.status, "missed": self.missed,
}
@staticmethod
def from_json(d: dict[str, Any]) -> "Corridor":
return Corridor(
lo=d.get("lo"), hi=d.get("hi"), count=d.get("count"),
truncated=bool(d.get("truncated")),
pages_done=int(d.get("pages_done") or 0),
status=d.get("status") or "pending", missed=int(d.get("missed") or 0),
)
def planned_pages(self, page_size: int, max_pages: int) -> int:
if not self.count:
return 1
return max(1, min(max_pages, math.ceil(self.count / page_size)))
@dataclass
class Plan:
base_url: str
target: int
batch_id: str
platform: str = "avito" # старые планы (до --platform) читаются как avito
corridors: list[Corridor] = field(default_factory=list)
created_at: str = ""
def save(self, path: Path) -> None:
path.write_text(
json.dumps(
{
"version": 1, "base_url": self.base_url, "target": self.target,
"batch_id": self.batch_id, "platform": self.platform,
"created_at": self.created_at,
"corridors": [c.to_json() for c in self.corridors],
},
ensure_ascii=False, indent=1,
),
encoding="utf-8",
)
@staticmethod
def load(path: Path) -> "Plan":
d = json.loads(path.read_text(encoding="utf-8"))
return Plan(
base_url=d["base_url"], target=int(d["target"]), batch_id=d["batch_id"],
platform=d.get("platform", "avito"),
created_at=d.get("created_at", ""),
corridors=[Corridor.from_json(c) for c in d.get("corridors", [])],
)
def _avito_build_url(base_url: str, page: int, lo: int | None, hi: int | None) -> str:
"""URL коридора Авито: pmin/pmax + пагинация. Гео-параметры не добавляем — #3043."""
parts = urlsplit(base_url)
q = [(k, v) for k, v in parse_qsl(parts.query, keep_blank_values=True)
if k not in {"p", "pmin", "pmax"}]
if lo is not None:
q.append(("pmin", str(int(lo))))
if hi is not None:
q.append(("pmax", str(int(hi))))
if page > 1:
q.append(("p", str(page)))
return urlunsplit(
(parts.scheme, parts.netloc, parts.path, urlencode(q), parts.fragment)
)
def _cian_build_url(base_url: str, page: int, lo: int | None, hi: int | None) -> str:
"""URL коридора Циан: minprice/maxprice + пагинация (аналог _avito_build_url).
Названия ключей — как в CianScraper._build_url (minprice/maxprice/p), но домен
и мульти-регион берём из base_url как есть, а не пересобираем через scraper.
"""
parts = urlsplit(base_url)
q = [(k, v) for k, v in parse_qsl(parts.query, keep_blank_values=True)
if k not in {"p", "minprice", "maxprice"}]
if lo is not None:
q.append(("minprice", str(int(lo))))
if hi is not None:
q.append(("maxprice", str(int(hi))))
if page > 1:
q.append(("p", str(page)))
return urlunsplit(
(parts.scheme, parts.netloc, parts.path, urlencode(q), parts.fragment)
)
def geometric_mid(lo: int | None, hi: int) -> int:
"""Геометрическая середина коридора.
Цены логнормальны: арифметическая середина 1 млн..100 млн (≈50 млн)
отрезает вырожденно-пустую верхнюю половину. sqrt(lo*hi) делит выборку
заметно ровнее.
"""
low = max(int(lo or PRICE_FLOOR), 1)
mid = int(math.sqrt(low * float(hi)))
return max(low + 1, min(hi - 1, mid))
def width_ratio(lo: int | None, hi: int | None) -> float:
if hi is None:
return float("inf")
return float(hi) / max(float(lo or PRICE_FLOOR), 1.0)
# --- платформенные адаптеры --------------------------------------------------
def _avito_detect_block(html: str) -> tuple[str, str] | None:
"""Доп. текстовые маркеры блока Авито поверх общих HTTP-проверок в _guard.
Порядок (firewall перед PoW) идентичен исходному _guard — регресс-нейтрально.
"""
if _is_firewall_page(html):
return "firewall", "firewall-страница на HTTP 200"
head = html[:4096].lower()
if any(m in head for m in _POW_MARKERS):
return "challenge", "PoW / проверка безопасности"
return None
async def _avito_wait_ready(page: Any) -> None:
"""Карточка Авито — обычный DOM-узел, ждём её появления."""
try:
await page.wait_for_selector('[data-marker="item"]', timeout=7_000)
except Exception: # noqa: BLE001 — пустая/блочная страница разбирается ниже
pass
async def _cian_wait_ready(page: Any) -> None:
"""У Циан карточки читаются не из DOM, а из window._cianConfig['frontend-serp'],
поэтому ждать `[data-marker="item"]` (авитовский селектор) бессмысленно: он не
появится никогда, и каждая загрузка стоила бы лишних 7 с таймаута. Ждём сам
state — ровно то, что потом парсится.
"""
try:
await page.wait_for_function("!!window._cianConfig", timeout=10_000)
except Exception: # noqa: BLE001 — пустая/блочная страница разбирается ниже
pass
def _cian_detect_block(html: str) -> tuple[str, str] | None:
"""Капча Циан: HTTP 200 и обычная с виду страница, но без выдачи.
Поймана живьём 10.09: тот же IP, что нормально отдаёт SERP браузеру и curl,
на голом http-клиенте получил 40 КБ с <title>Captcha - база объявлений ЦИАН</title>
и без window._cianConfig вообще. То есть блок Циан НЕ приходит ни 403, ни 429 —
только телом, и по HTTP-статусу его не отличить.
Маркер берём в первых 4 КБ: в блок-странице слово встречается там 4 раза
(title + текст), в нормальной выдаче на 2.6 МБ — ни разу (все восемь вхождений
лежат глубоко в скриптах). Ложных срабатываний на живой выдаче нет.
Вторая сеть на случай другой формы блока — parse_page: count is None и сырых
карточек 0 тоже даёт Blocked("challenge").
"""
if "captcha" in html[:4096].lower():
return "challenge", "капча Циан (HTTP 200, страница без выдачи)"
return None
def _cian_parse_cards(scraper: CianScraper, html: str) -> list[Any]:
"""Парс карточек Циан. Ничего не отбрасываем, только считаем ссылки на ЖК.
Кит фильтрует новостройки по `offer.newbuilding.id` (serp.py:401-402,
listing_segment == "novostroyki"), считая SERP-параметр object_type=1
ненадёжным. Для ЕКБ, где выдача бралась без него, это было верно, здесь —
нет: базовый URL уже несёт object_type[0]=1, и Циан подтверждает скоуп
сам (puid3=sale_type_second).
Замер живой страницы 10.09: из 28 карточек 16 имеют блок newbuilding, и у
ВСЕХ шестнадцати isFromBuilder=false, isFromLeadFactory=false, flatType=rooms,
а ЖК — «Пригород Лесное», «РУСИЧ Новые Котельники» и подобные. То есть это
обычная вторичка в новых домах, а не продажа от застройщика: newbuilding.id
означает «дом в ЖК», а не «лот застройщика».
С фильтром прогон терял бы 57% корпуса безвозвратно, а msk_raw — сырьё:
payload несёт listing_segment целиком, и отделить сегмент можно потом, на
импорте в listings. Поэтому здесь только счётчик last_nb_ref для notes.
"""
lots = scraper._parse_serp_html(html)
# Сырое число карточек — по нему parse_page отличает блок от пустой страницы.
scraper.last_raw_count = len(lots)
nb_ref = sum(1 for lot in lots
if getattr(lot, "listing_segment", None) == "novostroyki")
if nb_ref:
scraper.last_nb_ref = getattr(scraper, "last_nb_ref", 0) + nb_ref
return lots
@dataclass(frozen=True)
class PlatformAdapter:
"""Платформо-зависимые куски сбора. Всё общее (бисекция, guard по HTTP-статусу,
накопитель/заливка) — в основном теле скрипта и от платформы не зависит."""
name: str
table: str # msk_raw.<table>
default_base_url: str
page_size: int
max_pages: int
build_url: Callable[[str, int, int | None, int | None], str]
extract_total_count: Callable[[Any, str], int | None]
parse_cards: Callable[[Any, str], list[Any]]
detect_block: Callable[[str], tuple[str, str] | None]
wait_ready: Callable[[Any], Any] # корутина: дождаться готовности страницы
@property
def hard_cap(self) -> int:
return self.page_size * self.max_pages
def make_scraper(self) -> Any:
if self.name == "avito":
# avito_serp_ekb_only=False обязателен: с True парсер выбрасывает всё,
# где в URL нет /ekaterinburg/ — то есть все подмосковные слаги
# (serp.py:2154).
return AvitoScraper(
SimpleNamespace(avito_serp_ekb_only=False), # type: ignore[arg-type]
target_city_slug="moskva",
)
# glitchtip_dsn=None — отключает попытку sentry-репорта schema-regression
# из _report_schema_regression (нет реального DSN в локальном прогоне).
return CianScraper(SimpleNamespace(glitchtip_dsn=None)) # type: ignore[arg-type]
ADAPTERS: dict[str, PlatformAdapter] = {
"avito": PlatformAdapter(
name="avito",
table="avito_cards",
default_base_url=DEFAULT_AVITO_BASE_URL,
page_size=AVITO_PAGE_SIZE,
max_pages=AVITO_MAX_PAGES,
build_url=_avito_build_url,
extract_total_count=lambda scraper, html: scraper._extract_total_count(html),
parse_cards=lambda scraper, html: scraper._parse_html(html, "https://www.avito.ru"),
detect_block=_avito_detect_block,
wait_ready=_avito_wait_ready,
),
"cian": PlatformAdapter(
name="cian",
table="cian_cards",
default_base_url=DEFAULT_CIAN_BASE_URL,
page_size=CIAN_PAGE_SIZE,
max_pages=CIAN_MAX_PAGES,
build_url=_cian_build_url,
extract_total_count=lambda scraper, html: scraper._extract_total_offers(html),
parse_cards=_cian_parse_cards,
detect_block=_cian_detect_block,
wait_ready=_cian_wait_ready,
),
}
# --- загрузка страницы -----------------------------------------------------
def _guard(html: str, status: int | None, adapter: PlatformAdapter) -> None:
"""Порядок проверок фиксирован заданием; первое срабатывание = стоп.
HTTP-статус — платформо-независимая проверка. Текстовые маркеры (firewall/
PoW и т.п.) — через adapter.detect_block, у каждой платформы свои.
"""
if status in (403, 439):
raise Blocked("platform", f"HTTP {status}")
if status == 429:
raise Blocked("ratelimit", "HTTP 429")
extra = adapter.detect_block(html)
if extra is not None:
raise Blocked(*extra)
class Loader:
"""Одна СВОЯ вкладка в уже открытом Chrome владельца (CDP).
Ни браузер, ни контекст, ни чужие вкладки не закрываются и не трогаются:
это рабочий Chrome с залогиненным техаккаунтом.
"""
def __init__(self, delay: float, page_budget: int | None, adapter: PlatformAdapter) -> None:
self._delay = delay
self._budget = page_budget
self._adapter = adapter
self.loads = 0
self._page: Any = None
self._pw: Any = None
self._browser: Any = None
self._ctx: Any = None
self._last_load = 0.0
self._loads_on_page = 0
async def __aenter__(self) -> "Loader":
from playwright.async_api import async_playwright
from playwright.async_api import TimeoutError as _PwTimeout
global PlaywrightTimeoutError
PlaywrightTimeoutError = _PwTimeout
# AVITO_CDP — исторически названо под первую платформу, но это адрес
# браузера владельца, а не площадки: используется для обеих платформ.
endpoint = os.environ.get("AVITO_CDP", "http://localhost:9222")
self._pw = await async_playwright().start()
try:
self._browser = await self._pw.chromium.connect_over_cdp(endpoint)
except Exception as exc: # noqa: BLE001 — подсказка важнее типа
await self._pw.stop()
raise SystemExit(
f"Не удалось подключиться по CDP к {endpoint}: {exc}\n"
"Запусти Chrome с залогиненным техаккаунтом и ключом "
"--remote-debugging-port=9222, либо укажи адрес в AVITO_CDP."
) from exc
if not self._browser.contexts:
await self._pw.stop()
raise SystemExit(
"В подключённом Chrome нет ни одного контекста. Открой обычное окно "
"Chrome, запущенное с --remote-debugging-port=9222."
)
self._ctx = self._browser.contexts[0]
self._page = await self._ctx.new_page()
return self
async def __aexit__(self, *exc: object) -> None:
if self._page is not None:
try:
await self._page.close() # ТОЛЬКО своя вкладка
except Exception: # noqa: BLE001
pass
if self._pw is not None:
try:
await self._pw.stop()
except Exception: # noqa: BLE001
pass
def budget_left(self) -> bool:
return self._budget is None or self.loads < self._budget
async def _pause(self) -> None:
if self._last_load == 0.0:
return
jitter = self._delay * random.uniform(-0.2, 0.2)
wait = max(0.0, self._delay + jitter - (time.monotonic() - self._last_load))
if wait > 0:
print(f" пауза {wait:.1f} с", flush=True)
await asyncio.sleep(wait)
async def fetch(self, url: str) -> tuple[str, int | None]:
if not self.budget_left():
raise BudgetExhausted()
await self._recycle_if_needed()
await self._pause()
resp = await self._goto(url)
self._loads_on_page += 1
self.loads += 1
self._last_load = time.monotonic()
status = resp.status if resp is not None else None
await self._adapter.wait_ready(self._page)
html = await self._read_content()
_guard(html, status, self._adapter)
return html, status
async def _recycle_if_needed(self) -> None:
"""Каждые PAGE_RECYCLE_EVERY загрузок пересоздаём свою вкладку.
Рендерер Chrome накапливает память по всем навигациям вкладки, а проход
по Москве — под тысячу страниц в одной. Наблюдалось живьём: вкладка
падала с «Опаньки… Код ошибки: Out of Memory» при 34 ГБ свободных в
системе, то есть упирался именно рендерер, а не машина. Свежая вкладка
стоит одну навигацию и обнуляет счёт.
Закрывается ТОЛЬКО своя вкладка; контекст и чужие вкладки владельца не
трогаются — это его рабочий Chrome.
"""
if self._loads_on_page < PAGE_RECYCLE_EVERY:
return
print(f" вкладка пересоздаётся после {self._loads_on_page} загрузок "
"(память рендерера)", flush=True)
old = self._page
self._page = await self._ctx.new_page()
self._loads_on_page = 0
try:
await old.close()
except Exception: # noqa: BLE001 — старая вкладка могла уже умереть
pass
async def _goto(self, url: str, attempts: int = 3):
"""goto с ограниченным ретраем на таймаут навигации.
Авито изредка держит соединение до упора и goto падает по timeout. Это
НЕ признак отказа: в наблюдавшемся случае вкладка показывала нормальную
выдачу, а маркеров фаервола/PoW не было. Но и молча ретраить бесконечно
нельзя — тихий отказ выглядит ровно так же. Поэтому: перед каждым
повтором пробуем прочитать то, что в документе, и прогнать через _guard,
чтобы настоящий блок остановил прогон с правильной причиной; исчерпали
попытки — жёсткий стоп с причиной nav_timeout.
"""
for i in range(attempts):
try:
return await self._page.goto(url, wait_until="domcontentloaded",
timeout=90_000)
except PlaywrightTimeoutError:
try:
partial = await self._page.content()
except Exception: # noqa: BLE001 — документа может не быть вовсе
partial = ""
if partial:
_guard(partial, None, self._adapter) # настоящий блок остановит прогон здесь
if i == attempts - 1:
raise Blocked("nav_timeout") from None
print(f" таймаут навигации, попытка {i + 2}/{attempts}", flush=True)
await asyncio.sleep(10.0)
async def _read_content(self, attempts: int = 4) -> str:
"""page.content() с узким ретраем на гонку клиентской перенавигации.
Авито дорисовывает выдачу после domcontentloaded, и content() иногда
попадает ровно в момент смены документа: "Unable to retrieve content
because the page is navigating and changing the content". Это НЕ отказ
площадки — гвардов не касается, поэтому ретраим только эту ошибку и
только её, а любую другую поднимаем как есть.
"""
last: Exception | None = None
for i in range(attempts):
try:
return await self._page.content()
except Exception as exc: # noqa: BLE001 — сузили проверкой текста ниже
if "page is navigating" not in str(exc):
raise
last = exc
print(f" content() поймал перенавигацию, попытка {i + 2}/{attempts}",
flush=True)
await asyncio.sleep(1.5)
raise RuntimeError(f"page.content() не отдал документ за {attempts} попыток") from last
# --- заливка в msk_raw -----------------------------------------------------
def _sql_str(value: str) -> str:
return "'" + value.replace("'", "''") + "'"
def _csv_rows(rows: Iterable[dict[str, Any]]) -> str:
buf = io.StringIO()
writer = csv.writer(buf, lineterminator="\n")
for r in rows:
writer.writerow([
r["source_id"], r["observed_at"], r["batch_id"], r["kind"],
r["url"], r["price"], r["payload"],
])
return buf.getvalue()
def build_sql(batch_id: str, query: str, rows: list[dict[str, Any]],
started_at: str, table: str, kind: str = "serp") -> str:
"""Один поток на `psql -f -`: batch (FK!) → TEMP staging → \\copy → INSERT.
Одиночный `psql -c` через ssh ломается на квотинге скобок и кавычек, поэтому
только поток. rows_new = разница count(*) по batch_id до и после вставки.
table — msk_raw.avito_cards / msk_raw.cian_cards (из PlatformAdapter, не из
пользовательского ввода — подстановка f-строкой безопасна).
У staging обязателен INCLUDING IDENTITY: если целевая таблица объявлена как
GENERATED ALWAYS AS IDENTITY, один INCLUDING DEFAULTS identity не переносит,
а NOT NULL переносится всегда — и \\copy без колонки id падает на каждом
батче.
"""
bid = _sql_str(batch_id)
tbl = f"msk_raw.{table}"
return (
"BEGIN;\n"
"INSERT INTO msk_raw.batches (batch_id, kind, query, started_at)\n"
f"VALUES ({bid}, {_sql_str(kind)}, {_sql_str(query)}, "
f"CAST({_sql_str(started_at)} AS timestamptz))\n"
"ON CONFLICT (batch_id) DO NOTHING;\n"
f"CREATE TEMP TABLE _stg (LIKE {tbl} INCLUDING DEFAULTS INCLUDING IDENTITY) "
"ON COMMIT DROP;\n"
"CREATE TEMP TABLE _before ON COMMIT DROP AS\n"
f" SELECT count(*) AS n FROM {tbl} WHERE batch_id = {bid};\n"
"\\copy _stg (source_id,observed_at,batch_id,kind,url,price,payload) "
"FROM STDIN WITH (FORMAT csv)\n"
+ _csv_rows(rows)
+ "\\.\n"
f"INSERT INTO {tbl} "
"(source_id,observed_at,batch_id,kind,url,price,payload)\n"
"SELECT source_id,observed_at,batch_id,kind,url,price,payload FROM _stg\n"
"ON CONFLICT (source_id,batch_id,kind) DO NOTHING;\n"
"UPDATE msk_raw.batches b SET\n"
" rows_sent = coalesce(b.rows_sent,0) + (SELECT count(*) FROM _stg),\n"
" rows_new = coalesce(b.rows_new,0) +\n"
f" ((SELECT count(*) FROM {tbl} WHERE batch_id = {bid})\n"
" - (SELECT n FROM _before))\n"
f"WHERE b.batch_id = {bid};\n"
"COMMIT;\n"
)
def build_finalize_sql(batch_id: str, query: str, notes: str) -> str:
bid = _sql_str(batch_id)
return (
"INSERT INTO msk_raw.batches (batch_id, kind, query, started_at)\n"
f"VALUES ({bid}, 'serp', {_sql_str(query)}, now())\n"
"ON CONFLICT (batch_id) DO NOTHING;\n"
f"UPDATE msk_raw.batches SET finished_at = now(), notes = {_sql_str(notes)}\n"
f"WHERE batch_id = {bid};\n"
)
def run_psql(sql: str, ssh_host: str, container: str, db_user: str, db_name: str,
attempts: int = 4) -> None:
"""Заливка батча через ssh с ретраем на обрыв транспорта.
Прогон длится часами, и ssh рвётся: живьём поймано «Connection reset by peer»
(ssh возвращает 255) прямо посреди заливки — весь прогон умирал, а несброшенный
батч терялся. Ретраить безопасно: SQL идемпотентен (batch через ON CONFLICT DO
NOTHING, карточки через ON CONFLICT (source_id,batch_id,kind) DO NOTHING).
Ретраится ТОЛЬКО транспорт (ssh 255). Ошибка самого psql (ON_ERROR_STOP, любой
другой код) — это дефект данных или SQL, её повтор не лечит: поднимаем сразу.
"""
cmd = [
"ssh", ssh_host,
f"docker exec -i {container} psql -U {db_user} -d {db_name} "
"-v ON_ERROR_STOP=1 -f -",
]
for i in range(attempts):
proc = subprocess.run(cmd, input=sql.encode("utf-8"), capture_output=True)
out = (proc.stdout + proc.stderr).decode("utf-8", "replace").strip()
if proc.returncode == 0:
if out:
print(f" psql: {out}", flush=True)
return
if proc.returncode != 255 or i == attempts - 1:
raise RuntimeError(f"psql через ssh вернул {proc.returncode}:\n{out}")
tail = out.splitlines()[-1] if out else "без вывода"
print(f" ssh оборвался ({tail}), повтор заливки {i + 2}/{attempts}", flush=True)
time.sleep(15.0 * (i + 1))
def psql_scalar(sql: str, ssh_host: str, container: str, db_user: str,
db_name: str) -> str:
"""Тот же ssh+psql, что и run_psql, но `-Atc` и с возвратом stdout.
Нужен для дешёвых проверок до старта (преflight целевой таблицы в collect):
один SELECT, без ретраев и без потока на stdin.
"""
cmd = [
"ssh", ssh_host,
f"docker exec -i {container} psql -U {db_user} -d {db_name} "
f"-v ON_ERROR_STOP=1 -Atc {shlex.quote(sql)}",
]
proc = subprocess.run(cmd, capture_output=True)
if proc.returncode != 0:
err = (proc.stdout + proc.stderr).decode("utf-8", "replace").strip()
raise RuntimeError(f"psql через ssh вернул {proc.returncode}:\n{err}")
return proc.stdout.decode("utf-8", "replace").strip()
# --- накопитель карточек ---------------------------------------------------
@dataclass
class Sink:
"""Батчами на прод (ssh+psql) или в локальный CSV при --dry-run."""
batch_id: str
started_at: str
query: str
table: str
batch_size: int
dry_run: bool
csv_path: Path
ssh_host: str
container: str
db_user: str
db_name: str
buffer: list[dict[str, Any]] = field(default_factory=list)
sent: int = 0
skipped_non_numeric: int = 0
def add(self, lot: Any) -> None:
raw_id = str(getattr(lot, "source_id", "") or "")
try:
source_id = int(raw_id) # в БД bigint, у ScrapedLot — строка
except (TypeError, ValueError):
self.skipped_non_numeric += 1
return
payload = lot.model_dump(mode="json")
self.buffer.append({
"source_id": source_id,
"observed_at": datetime.now(timezone.utc).isoformat(),
"batch_id": self.batch_id,
"kind": "serp",
"url": payload.get("source_url"),
"price": payload.get("price_rub"),
"payload": json.dumps(payload, ensure_ascii=False),
})
def maybe_flush(self) -> None:
if len(self.buffer) >= self.batch_size:
self.flush()
def flush(self) -> None:
if not self.buffer:
return
rows, self.buffer = self.buffer, []
if self.dry_run:
fresh = not self.csv_path.exists()
with self.csv_path.open("a", encoding="utf-8", newline="") as fh:
if fresh:
fh.write("source_id,observed_at,batch_id,kind,url,price,payload\n")
fh.write(_csv_rows(rows))
print(f" [dry-run] {len(rows)} строк → {self.csv_path}", flush=True)
else:
run_psql(
build_sql(self.batch_id, self.query, rows, self.started_at, self.table),
self.ssh_host, self.container, self.db_user, self.db_name,
)
print(f" залито {len(rows)} строк в msk_raw.{self.table}", flush=True)
self.sent += len(rows)
def finalize(self, notes: str) -> None:
self.flush()
if self.dry_run:
print(f" [dry-run] finalize: {notes}", flush=True)
return
run_psql(
build_finalize_sql(self.batch_id, self.query, notes),
self.ssh_host, self.container, self.db_user, self.db_name,
)
# --- сбор ------------------------------------------------------------------
def parse_page(scraper: Any, html: str, url: str,
adapter: PlatformAdapter) -> tuple[int | None, list[Any]]:
"""Счётчик и карточки страницы + детект блока по СЫРОМУ числу карточек.
Гвард считает по сырому числу (до отброса новостроек), а не по итоговому
списку: на Циан страница может легально дать 0 карточек после фильтра
вторички, и такой штатный ноль неотличим от блока. scraper.last_raw_count
выставляет _cian_parse_cards; у Авито этого атрибута нет,
поэтому дефолт getattr — len(lots), и авито-ветка ведёт себя как раньше.
Это вторая сеть после текстовых маркеров в adapter.detect_block: блок,
не попавший под маркер, всплывает как проваленный extract_state, то есть
count=None и пустой парс одновременно.
"""
count = adapter.extract_total_count(scraper, html)
lots = adapter.parse_cards(scraper, html)
raw = getattr(scraper, "last_raw_count", len(lots))
if count is None and raw == 0:
raise Blocked(
"challenge",
f"страница без счётчика и без карточек (провал extract_state): {url}",
)
if raw == 0 and count:
raise Blocked("empty_page", f"0 карточек при счётчике {count}: {url}")
return count, lots
async def probe(loader: Loader, scraper: Any, base_url: str,
lo: int | None, hi: int | None,
adapter: PlatformAdapter) -> tuple[int | None, list[Any]]:
url = adapter.build_url(base_url, 1, lo, hi)
html, _ = await loader.fetch(url)
return parse_page(scraper, html, url, adapter)
async def build_plan(loader: Loader, scraper: Any, base_url: str, target: int,
cache: dict[tuple[int | None, int | None], list[Any]],
adapter: PlatformAdapter,
) -> list[Corridor]:
"""Адаптивная бисекция по цене; страница 1 каждого коридора кэшируется."""
corridors: list[Corridor] = []
hard_cap = adapter.hard_cap
def emit(lo: int | None, hi: int | None, count: int | None,
truncated: bool, lots: list[Any]) -> None:
missed = max(0, (count or 0) - hard_cap) if truncated else 0
c = Corridor(lo=lo, hi=hi, count=count, truncated=truncated, missed=missed)
corridors.append(c)
cache[(lo, hi)] = lots
flag = " TRUNCATED" if truncated else ""
print(f" коридор {c.label()} count={count} "
f"страниц={c.planned_pages(adapter.page_size, adapter.max_pages)}{flag}",
flush=True)
if truncated:
print(f" ВНИМАНИЕ: коридор {c.label()} не влезает в потолок "
f"{hard_cap}; заведомо не добрано ~{missed} объявлений", flush=True)
async def find_upper(lo: int | None) -> int:
"""Верхнюю границу открытого коридора ищем удвоением от разумного старта."""
cand = max(int(lo or PRICE_FLOOR) * 2, PRICE_PROBE_START)
while cand < PRICE_CEIL:
cnt, _ = await probe(loader, scraper, base_url, cand, None, adapter)
print(f" проба хвоста pmin={cand:_} count={cnt}", flush=True)
if cnt is not None and cnt <= target:
return cand
cand *= 2
return cand
async def split(lo: int | None, hi: int | None, depth: int,
count: int | None, lots: list[Any]) -> None:
if count is None:
url = adapter.build_url(base_url, 1, lo, hi)
raise Blocked("empty_page", f"счётчик не прочитался: {url}")
if count <= target:
emit(lo, hi, count, False, lots)
return
if depth >= MAX_DEPTH or width_ratio(lo, hi) <= MIN_WIDTH_RATIO:
# Предохранитель: не молчим — помечаем truncated и считаем недобор.
emit(lo, hi, count, count > hard_cap, lots)
return
upper = hi if hi is not None else await find_upper(lo)
if hi is None:
tail_cnt, tail_lots = await probe(loader, scraper, base_url, upper, None, adapter)
emit(upper, None, tail_cnt,
bool(tail_cnt and tail_cnt > hard_cap), tail_lots)
mid = geometric_mid(lo, upper)
for sub_lo, sub_hi in ((lo, mid), (mid, upper)):
sub_cnt, sub_lots = await probe(loader, scraper, base_url, sub_lo, sub_hi, adapter)
print(f" проба {sub_lo or '-'}..{sub_hi} count={sub_cnt}", flush=True)
await split(sub_lo, sub_hi, depth + 1, sub_cnt, sub_lots)
root_cnt, root_lots = await probe(loader, scraper, base_url, None, None, adapter)
print(f"Всего по базовому запросу: {root_cnt}", flush=True)
await split(None, None, 0, root_cnt, root_lots)
return corridors
async def collect(args: argparse.Namespace) -> int:
adapter = ADAPTERS[args.platform]
scraper = adapter.make_scraper()
out_dir = Path(args.out_dir).resolve()
out_dir.mkdir(parents=True, exist_ok=True)
plan_path = out_dir / f"plan-{args.batch_id}.json"
csv_path = out_dir / f"cards-{args.batch_id}.csv"
started_at = datetime.now(timezone.utc).isoformat()
plan: Plan | None = None
if args.resume:
if not plan_path.exists():
print(f"--resume: плана нет — {plan_path}", file=sys.stderr)
return 1
plan = Plan.load(plan_path)
done = sum(1 for c in plan.corridors if c.status == "done")
print(f"Resume по {plan_path}: коридоров {len(plan.corridors)}, "
f"готово {done}", flush=True)
# Resume: URL и платформу берём из сохранённого плана, а не из CLI — коридоры
# посчитаны именно под них. Расхождение = молчаливая заливка чужой выдачи под
# тем же batch_id, поэтому это ошибка, а не тихий приоритет одного из двух.
if plan is not None and plan.platform != args.platform:
raise SystemExit(
"--resume: план построен для другой платформы."
f" В плане {plan.platform}, в аргументах {args.platform}."
)
if plan is not None and plan.base_url != args.base_url:
raise SystemExit(
"--resume: план построен для другого URL."
f" В плане {plan.base_url}, в аргументах {args.base_url}."
" Убери --base-url (возьмётся из плана) либо начни новый batch_id."
)
base_url = plan.base_url if plan is not None else args.base_url
page_budget = None if args.full else args.measure
mode = "FULL" if args.full else f"MEASURE<={page_budget}"
print(f"Режим: {mode}; platform={args.platform}; batch_id={args.batch_id}; "
f"delay={args.delay}s; target={args.target_count}; dry_run={args.dry_run}",
flush=True)
# Преflight целевой таблицы: без неё прогон умирал бы только на первой
# заливке — после часов планирования и сбора. Проверяем на первой секунде.
if not args.dry_run:
if not psql_scalar(f"SELECT to_regclass('msk_raw.{adapter.table}')",
args.ssh_host, args.container,
args.db_user, args.db_name):
raise SystemExit(
f"Таблицы msk_raw.{adapter.table} на проде нет. Применить"
" tradein-mvp/backend/data/sql/"
"299_msk_raw_cian_domclick_yandex_cards.sql"
" либо гонять с --dry-run."
)
sink = Sink(
batch_id=args.batch_id, started_at=started_at, query=base_url,
table=adapter.table, batch_size=args.batch_size, dry_run=args.dry_run,
csv_path=csv_path, ssh_host=args.ssh_host, container=args.container,
db_user=args.db_user, db_name=args.db_name,
)
cache: dict[tuple[int | None, int | None], list[Any]] = {}
total = 0
stop_reason = ""
rc = 0
loads = 0
async with Loader(args.delay, page_budget, adapter) as loader:
try:
if plan is None:
print("Строю план коридоров...", flush=True)
corridors = await build_plan(loader, scraper, base_url,
args.target_count, cache, adapter)
plan = Plan(base_url=base_url, target=args.target_count,
batch_id=args.batch_id, platform=args.platform,
corridors=corridors, created_at=started_at)
plan.save(plan_path)
print(f"План сохранён: {plan_path} ({len(corridors)} коридоров)",
flush=True)
for corridor in plan.corridors:
if corridor.status == "done":
continue
pages = corridor.planned_pages(adapter.page_size, adapter.max_pages)
print(f"Коридор {corridor.label()} count={corridor.count} "
f"страниц={pages} (с {corridor.pages_done + 1})", flush=True)
for page in range(corridor.pages_done + 1, pages + 1):
key = (corridor.lo, corridor.hi)
if page == 1 and key in cache:
lots = cache.pop(key) # страница 1 уже скачана при планировании
else:
url = adapter.build_url(base_url, page, corridor.lo, corridor.hi)
html, _ = await loader.fetch(url)
_, lots = parse_page(scraper, html, url, adapter)
for lot in lots:
sink.add(lot)
total += len(lots)
corridor.pages_done = page
print(f" стр.{page}/{pages}: карточек {len(lots)}, "
f"итого {total}", flush=True)
sink.maybe_flush()
plan.save(plan_path)
if not lots:
print(" пустая страница — конец коридора", flush=True)
break
corridor.status = "done"
plan.save(plan_path)
except BudgetExhausted:
stop_reason = "потолок --measure исчерпан"
print(f"Стоп: {stop_reason}", flush=True)
except Blocked as exc:
stop_reason = f"BLOCKED/{exc.reason}: {exc.detail}"
print(f"СТОП: {stop_reason}", file=sys.stderr, flush=True)
rc = 2
finally:
loads = loader.loads
if plan is not None:
plan.save(plan_path)
truncated = [c for c in (plan.corridors if plan else []) if c.truncated]
missed = sum(c.missed for c in truncated)
nb_ref = getattr(scraper, "last_nb_ref", 0)
notes = "; ".join(x for x in [
f"mode={mode}", f"platform={args.platform}", f"loads={loads}", f"cards={total}",
f"skipped_non_numeric={sink.skipped_non_numeric}",
(f"nb_ref={nb_ref}" if nb_ref else ""),
(f"truncated_corridors={len(truncated)} missed~{missed}" if truncated else ""),
stop_reason,
] if x)
try:
sink.finalize(notes)
except Exception as exc: # noqa: BLE001 — не прятать исходную причину стопа
print(f"finalize провалился: {exc}", file=sys.stderr)
rc = rc or 1
print(f"Готово. Загрузок: {loads}; карточек: {total}; отправлено: {sink.sent}; "
f"пропущено нечисловых source_id: {sink.skipped_non_numeric}; "
f"notes: {notes}", flush=True)
return rc
def parse_args(argv: list[str] | None = None) -> argparse.Namespace:
p = argparse.ArgumentParser(
prog="collect.py",
description="Ручной сбор SERP Авито/Циан (вторичка, Москва+МО) в прод-схему msk_raw.",
)
p.add_argument("--platform", choices=tuple(ADAPTERS), default="avito",
help="площадка сбора (дефолт avito)")
p.add_argument("--base-url", default=None,
help="базовый URL выдачи (дефолт зависит от --platform)")
p.add_argument("--measure", type=int, default=100, metavar="N",
help="режим замера: не больше N загрузок страниц (дефолт 100)")
p.add_argument("--full", action="store_true",
help="полный проход без потолка страниц (включается только явно)")
p.add_argument("--dry-run", action="store_true",
help="ничего не слать на прод, писать CSV локально")
p.add_argument("--resume", action="store_true",
help="продолжить по сохранённому плану коридоров")
p.add_argument("--delay", type=float, default=8.0,
help="пауза между загрузками, с (±20%% джиттер, дефолт 8.0)")
p.add_argument("--batch-size", type=int, default=1000,
help="карточек в одной заливке (дефолт 1000)")
p.add_argument("--target-count", type=int, default=1500,
help="целевой размер коридора; больше — делим (дефолт 1500)")
p.add_argument("--batch-id", default=None,
help="batch_id в msk_raw.batches (дефолт msk-serp-<platform>-<UTC>)")
p.add_argument("--out-dir", default=str(Path(__file__).resolve().parent / "runs"),
help="каталог плана/CSV")
p.add_argument("--ssh-host", default="selectel", help="ssh-хост прода")
p.add_argument("--container", default="tradein-postgres",
help="имя контейнера Postgres на проде")
p.add_argument("--db-user", default="tradein")
p.add_argument("--db-name", default="tradein")
args = p.parse_args(argv)
if args.base_url is None:
args.base_url = ADAPTERS[args.platform].default_base_url
if args.batch_id is None:
args.batch_id = (
f"msk-serp-{args.platform}-" + datetime.now(timezone.utc).strftime("%Y%m%d-%H%M%S")
)
if not _BATCH_ID_RE.match(args.batch_id):
p.error("--batch-id: допустимы только символы [A-Za-z0-9._-]")
if args.measure < 1:
p.error("--measure должен быть >= 1")
return args
def main(argv: list[str] | None = None) -> int:
return asyncio.run(collect(parse_args(argv)))
if __name__ == "__main__":
raise SystemExit(main())