gendesign/tradein-mvp/scripts/local-avito-msk/collect.py
lekss361 628f59dcc0
All checks were successful
Deploy Trade-In / test (push) Has been skipped
Deploy Trade-In / build-backend (push) Has been skipped
Deploy Trade-In / deploy (push) Successful in 1m1s
Deploy Trade-In / changes (push) Successful in 14s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / deploy-status (push) Successful in 2s
Deploy Trade-In / perimeter-smoke (push) Successful in 1m41s
ДомКлик — четвёртая площадка ручного сбора Москвы (#3496)
2026-09-12 12:05:18 +00:00

1574 lines
81 KiB
Python
Raw Blame History

This file contains invisible Unicode characters

This file contains invisible Unicode characters that are indistinguishable to humans but may be processed differently by a computer. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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 Авито/Циан/Яндекса/DomClick по Москве и МО (эпик #2989, трек 1).
Запускается ВРУЧНУЮ с машины владельца. Прод-скрейпер, его расписания и
прокси-пул не задействованы вообще: браузер — уже открытый Chrome владельца
(подключение по CDP), парсер — импорт из scraper-kit, заливка — поток в psql
через ssh. Скрипт ничего не устанавливает и своего профиля не поднимает.
Платформа выбирается ключом --platform {avito,cian,yandex} (дефолт 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,
)
# Яндекс ходит в gate-API и отдаёт JSON, а не SERP-разметку: DOM-парсера у него нет
# и тащить его сюда нечего. Из кита берём ровно чистые функции разбора gate-payload —
# класс YandexRealtyScraper не нужен (он существует ради BrowserFetcher/camoufox и
# пула прокси, а здесь транспорт — вкладка Chrome владельца).
from scraper_kit.providers.yandex.serp import ( # noqa: E402
_extract_gate_data,
_extract_json_from_content,
_is_gate_error,
_parse_gate_json,
)
# DomClick — тоже JSON BFF-ответ (Layer A прод-скрейпера), а не HTML SERP, как
# и Яндекс выше. Берём из кита чистую функцию извлечения JSON (_extract_json)
# и сам класс DomClickScraper — не ради транспорта (тот же приём с браузером
# владельца), а ради DomClickScraper._map_item: маппинг offer-item → ScrapedLot
# дублировать своим кодом было бы регрессивно. Маркеры блока — единый источник
# Layer A/B кита, см. докстринг DOMCLICK_BLOCK_MARKERS.
from scraper_kit.domclick_exceptions import ( # noqa: E402
DOMCLICK_BLOCK_MARKERS,
DomClickBlockedError,
)
from scraper_kit.providers.domclick.serp import ( # noqa: E402
DomClickScraper,
_extract_json,
)
# Вкладка, открытая у владельца: вторичка, Москва + МО.
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"
)
# Яндекс.Недвижимость: gate-API (JSON), а не HTML-выдача.
#
# rgid установлен ЭМПИРИЧЕСКИ 11.09, а не взят из памяти: в SSR-разметке
# realty.yandex.ru/moskva/kupit/kvartira/vtorichniy-rynok/ лежит
# "geo":{"id":213,"type":"CITY","rgid":587795,...,"name":"Москва",
# "parents":[{"id":1,"name":"Москва и МО","rgid":"741964",...}]}
# и там же page.params = {"rgid":587795,"type":"SELL","category":"APARTMENT",
# "newFlat":"NO"} — ровно параметры gate-запроса. Аналогично со страницы МО снят
# rgid 587654 ("Московская область", SUBJECT_FEDERATION).
#
# Счётчики вторички gate-API (pager.totalItems, живой запрос 11.09):
# rgid=587795 (Москва) 18 705
# rgid=587654 (МО) 14 242
# rgid=741964 (Москва и МО) 31 074
# rgid=559132 (ЕКБ, _EKB_RGID) 4 060 ← контроль: московский скоуп в 7.6 раза больше
# Дефолт — 741964: единый скоуп «Москва и МО», как region=-1 у Циан и
# moskva_i_mo у Авито.
_YANDEX_MSK_RGID = 587795
_YANDEX_MO_RGID = 587654
_YANDEX_MSK_MO_RGID = 741964
DEFAULT_YANDEX_BASE_URL = (
"https://realty.yandex.ru/gate/react-page/get/"
f"?rgid={_YANDEX_MSK_MO_RGID}&type=SELL&category=APARTMENT&newFlat=NO"
"&_pageType=search&_providers=react-search-results-data"
)
# Замерено живым проходом (не из документации): выдача Москва+МО отдаёт 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
# Замерено живым gate-запросом 11.09, не взято из кита: pager отдаёт
# pageSize=20 и totalPages=25 ДАЖЕ когда totalItems=31 074 (то есть 1554
# страницы по построению). page=26 отвечает HTTP 301 — редирект, не выдача.
# Значит жёсткий потолок Яндекса = 25*20 = 500 офферов на один набор фильтров,
# втрое ниже авитовского и цианского. Константа _GATE_MAX_PAGES_CAP=50 в
# scraper_kit — исторический запас, живьём недостижима.
YANDEX_PAGE_SIZE = 20
YANDEX_MAX_PAGES = 25
# Целевой размер коридора не может превышать hard_cap площадки, иначе бисекция
# штатно «дорезает» до 1500 и каждый коридор уезжает в truncated.
YANDEX_TARGET_COUNT = YANDEX_PAGE_SIZE * YANDEX_MAX_PAGES # 500
# DomClick: JSON BFF listing API, Москва (эпик #2989, трек 1).
#
# GUID Москвы проверен живым запросом 12.09: и region, и locality — один и тот
# же 1d1463ae-c80f-4d19-9331-a1b68a85b553 (для сравнения ЕКБ —
# 0d475b79-88de-4054-818c-37d8f9d0d440). Параметр aids (у ЕКБ 20561, сужает
# выдачу до конкретного агрегатора) для Москвы НЕ нужен — без него счётчик
# листинга точно совпадает с сайтом.
MSK_DOMCLICK_GUID = "1d1463ae-c80f-4d19-9331-a1b68a85b553"
DEFAULT_DOMCLICK_BASE_URL = (
"https://bff-search-web.domclick.ru/api/offers/v1"
f"?address={MSK_DOMCLICK_GUID}&deal_type=sale&category=living&offer_type=flat"
"&sort=qi&sort_dir=desc&limit=20&offset=0"
)
# Замерено живым запросом 12.09, не из документации: offset=1980 отдаёт полную
# страницу (20 items), offset=2000 отвечает HTTP 400
# {"statusCode":400,"error":"Bad Request"}. Значит жёсткий потолок пагинации —
# 100 страниц по 20 штук = 2000 офферов на один набор фильтров.
DOMCLICK_PAGE_SIZE = 20
DOMCLICK_MAX_PAGES = 100
# bbox Москвы с ТиНАО — единственный надёжный гео-гард для DomClick.
# offerRegionName использовать НЕЛЬЗЯ: часть офферов Новой Москвы приходит с
# именами вида "г. Говорово", а не "Москва" — гард по имени региона молча
# вырезал бы легитимные лоты. scraper._is_geo_ok() кита сюда тоже не подходит:
# он захардкожен на offerRegionName == "Екатеринбург" и отбросил бы буквально
# всю московскую выдачу.
_MSK_LAT_MIN, _MSK_LAT_MAX = 55.14, 56.02
_MSK_LON_MIN, _MSK_LON_MAX = 36.80, 37.97
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
PlaywrightError: type[BaseException] = Exception
# Сетевые отказы Chromium, которые НЕ означают отказ площадки: у машины моргнула
# сеть, сменился интерфейс, оборвалось соединение. Замер 12.09: полный проход
# Яндекса умер на восьмом часу на `net::ERR_NETWORK_CHANGED`, потеряв живой
# прогон целиком — площадка при этом отвечала 200 сразу после. Такие ошибки
# лечатся ожиданием, а не остановкой; всё остальное по-прежнему поднимается.
_TRANSIENT_NET_ERRORS = (
"net::ERR_NETWORK_CHANGED",
"net::ERR_INTERNET_DISCONNECTED",
"net::ERR_NAME_NOT_RESOLVED",
"net::ERR_CONNECTION_RESET",
"net::ERR_CONNECTION_CLOSED",
"net::ERR_CONNECTION_ABORTED",
"net::ERR_ADDRESS_UNREACHABLE",
"net::ERR_NETWORK_IO_SUSPENDED",
)
def _is_transient_net_error(exc: BaseException) -> bool:
text = str(exc)
return any(marker in text for marker in _TRANSIENT_NET_ERRORS)
# Через столько загрузок вкладка сборщика пересоздаётся (см. _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
# --- Яндекс: gate-API (JSON), общего с DOM-платформами только конвейер ---------
# Chrome отдаёт JSON-документ как текстовый узел внутри <pre>, а сериализатор
# экранирует в тексте ровно три символа. Разэкранируем их обратно ДО json.loads:
# иначе описания приезжают с &amp;/&lt; вместо & и <. Порядок обязателен —
# "&amp;" последним, иначе "&amp;lt;" из описания превратился бы в "<".
def _unescape_pre_text(text: str) -> str:
return text.replace("&lt;", "<").replace("&gt;", ">").replace("&amp;", "&")
def _yandex_payload(scraper: Any, html: str) -> dict[str, Any] | None:
"""gate-payload из содержимого вкладки, с мемоизацией на один HTML.
parse_page дёргает сначала extract_total_count, потом parse_cards — оба по
одной и той же строке. JSON тут на сотни КБ, поэтому разбираем один раз и
кладём на scraper вместе с самой строкой. Сравнение — `is`, и ссылка на
строку хранится: по id() было бы неверно, id освобождённой строки может
достаться следующей.
"""
if getattr(scraper, "_payload_html", None) is html:
return scraper._payload
json_text = _extract_json_from_content(html)
payload: dict[str, Any] | None = None
if json_text:
try:
payload = json.loads(_unescape_pre_text(json_text))
except (json.JSONDecodeError, ValueError):
payload = None
if payload is not None and _is_gate_error(payload):
payload = None
scraper._payload_html = html
scraper._payload = payload
return payload
def _yandex_build_url(base_url: str, page: int, lo: int | None, hi: int | None) -> str:
"""URL коридора Яндекса: priceMin/priceMax + page (1-based).
page=0 gate-API считает ошибкой, поэтому page проставляется ВСЕГДА, в том
числе на первой странице, — в отличие от Авито/Циан, где `p` на первой
странице опускается.
"""
parts = urlsplit(base_url)
q = [(k, v) for k, v in parse_qsl(parts.query, keep_blank_values=True)
if k not in {"page", "priceMin", "priceMax"}]
if lo is not None:
q.append(("priceMin", str(int(lo))))
if hi is not None:
q.append(("priceMax", str(int(hi))))
q.append(("page", str(max(1, page))))
return urlunsplit(
(parts.scheme, parts.netloc, parts.path, urlencode(q), parts.fragment)
)
def _yandex_extract_total(scraper: Any, html: str) -> int | None:
payload = _yandex_payload(scraper, html)
if payload is None:
return None
extracted = _extract_gate_data(payload)
if extracted is None:
return None
_entities, pager = extracted
total = pager.get("totalItems")
return int(total) if total is not None else None
def _yandex_parse_cards(scraper: Any, html: str) -> list[Any]:
"""Разбор gate-payload в ScrapedLot теми же функциями, что и прод-скрейпер.
page_param берём из самого ответа (response.pageParams.params.page), а не из
аргумента: контракт адаптера номера страницы не передаёт, а в raw_payload
должна лечь та страница, которую реально отдал gate.
"""
payload = _yandex_payload(scraper, html)
if payload is None:
scraper.last_raw_count = 0
return []
extracted = _extract_gate_data(payload)
if extracted is None:
scraper.last_raw_count = 0
return []
entities, _pager = extracted
# Сырое число — до отбраковки лотов без offerId/цены в _entity_to_lot.
scraper.last_raw_count = len(entities)
try:
page_param = int(payload["response"]["pageParams"]["params"]["page"])
except (KeyError, TypeError, ValueError):
page_param = 1
return _parse_gate_json(payload, page_param=page_param, new_flat="NO")
def _yandex_detect_block(html: str) -> tuple[str, str] | None:
"""Блок Яндекса: SmartCaptcha приходит HTML-страницей на месте JSON.
Порядок важен: сперва ищем маркеры капчи (иначе «нет JSON» рапортовалось бы
как schema-drift), потом проверяем, что документ вообще похож на gate-ответ.
Пустой gate-payload без response — тоже стоп: это заглушка, а не выдача.
"""
head = html[:4096].lower()
if "smartcaptcha" in head or "showcaptcha" in head or "captcha" in head:
return "challenge", "SmartCaptcha Яндекса вместо gate-ответа"
json_text = _extract_json_from_content(html)
if not json_text:
return "challenge", "во вкладке нет JSON (gate отдал не тот документ)"
return None
async def _yandex_wait_ready(page: Any) -> None:
"""gate-API — документ JSON, клиентской дорисовки нет.
Ждать `[data-marker=item]` (Авито) или window._cianConfig (Циан) здесь
бессмысленно — их не будет никогда, каждая загрузка стоила бы полного
таймаута. Chrome заворачивает текстовый документ в <pre>, его и ждём, коротко:
к моменту domcontentloaded он, как правило, уже на месте, а на капче его нет
вовсе — и тогда через 3 с отработает _yandex_detect_block.
"""
try:
await page.wait_for_selector("pre", timeout=3_000)
except Exception: # noqa: BLE001 — капча/заглушка разбирается detect_block ниже
pass
# --- DomClick: JSON BFF (конвейер как у Яндекса — payload вместо DOM) ---------
def _domclick_build_url(base_url: str, page: int, lo: int | None, hi: int | None) -> str:
"""URL коридора DomClick: sale_price__gte/__lte + offset (аналог _cian_build_url).
page — 1-based, как у остальных платформ; offset у BFF — 0-based, поэтому
offset = (page - 1) * DOMCLICK_PAGE_SIZE. offset проставляется ВСЕГДА, в
т.ч. на первой странице (как priceMin/priceMax/page у Яндекса): старые
offset/sale_price__gte/sale_price__lte из base_url выкидываются и ставятся
заново, дублей ключей не остаётся.
"""
parts = urlsplit(base_url)
q = [(k, v) for k, v in parse_qsl(parts.query, keep_blank_values=True)
if k not in {"offset", "sale_price__gte", "sale_price__lte"}]
if lo is not None:
q.append(("sale_price__gte", str(int(lo))))
if hi is not None:
q.append(("sale_price__lte", str(int(hi))))
q.append(("offset", str((page - 1) * DOMCLICK_PAGE_SIZE)))
return urlunsplit(
(parts.scheme, parts.netloc, parts.path, urlencode(q), parts.fragment)
)
def _unescape_pre_entities(html: str) -> str:
"""Снять HTML-экранирование текстового узла <pre>, в который Chrome заворачивает
ответ ручки.
`page.content()` отдаёт СЕРИАЛИЗОВАННЫЙ документ, а не тело ответа: сериализатор
экранирует в текстовом узле ровно четыре вещи — `&`, `<`, `>` и неразрывный
пробел. Синтаксис JSON от этого не страдает (`&amp;` внутри строки — валидная
строка), поэтому дефект тихий: `json.loads` проходит, а в описании и адресе
вместо «&» и неразрывного пробела оседают `&amp;` и `&nbsp;`. Замер на живой
выдаче Москвы: 29 `&nbsp;` на одной странице из двадцати карточек.
Кавычки в список не входят намеренно: сериализатор экранирует их только в
значениях атрибутов, а в текстовом узле `"` остаётся собой — и это важно, иначе
`\"` внутри JSON-строки превратился бы в невалидный escape и уронил разбор
страницы целиком.
Порядок фиксирован: `&amp;` разворачивается ПОСЛЕДНИМ, иначе уже развёрнутый
амперсанд склеится со следующей последовательностью и даст второй разбор.
"""
return (
html.replace("&lt;", "<")
.replace("&gt;", ">")
.replace("&nbsp;", " ")
.replace("&amp;", "&")
)
def _domclick_payload(scraper: Any, html: str) -> dict[str, Any] | None:
"""JSON BFF-ответ из содержимого вкладки, с мемоизацией на один HTML.
parse_page дёргает extract_total_count и parse_cards по одному и тому же
html — разбираем один раз (тот же приём, что и _yandex_payload). Парсер не
свой: _extract_json — та же функция, которой пользуется прод-скрейпер
(providers/domclick/serp.py), со своим корректным порядком «сначала JSON,
потом маркеры блока» (#3267 — маркер может встретиться в тексте самого
объявления, подстрочный поиск ДО попытки распарсить JSON уже забанил живой
узел на 6 часов в проде 30.08).
К моменту вызова этой функции _guard уже прогнал тот же html через
_domclick_detect_block, и настоящий блок остановил бы сбор раньше —
DomClickBlockedError/ValueError здесь только тихо превращаются в None, как
и в yandex-ветке.
"""
if getattr(scraper, "_payload_html", None) is html:
return scraper._payload
try:
payload: dict[str, Any] | None = _extract_json(_unescape_pre_entities(html))
except (DomClickBlockedError, ValueError):
payload = None
scraper._payload_html = html
scraper._payload = payload
return payload
def _domclick_extract_total(scraper: Any, html: str) -> int | None:
"""result.pagination.total — а НЕ отдельная count-ручка.
Замер живым запросом 12.09: pagination.total по Москве без фильтра комнат
= 23692, и это число реально достижимо постраничным обходом (offset до
1980, см. DOMCLICK_MAX_PAGES). У отдельной count-ручки другое число —
25898: оно включает дубли одного объекта у разных агентств, которые
пагинацией физически недостижимы. Если бы бисекция целилась в
count-ручку, она считала бы коридоры «недобранными» там, где добирать
нечего, и truncated/missed в notes врали бы.
"""
payload = _domclick_payload(scraper, html)
if not payload:
return None
result = payload.get("result")
if not isinstance(result, dict):
return None
pagination = result.get("pagination")
if not isinstance(pagination, dict):
return None
total = pagination.get("total")
return int(total) if total is not None else None
def _domclick_parse_cards(scraper: Any, html: str) -> list[Any]:
"""result.items → ScrapedLot через DomClickScraper._map_item (кит, без своего маппинга).
Гео-гард — СВОЙ, по bbox Москвы+ТиНАО, а не scraper._is_geo_ok(): тот метод
кита захардкожен на offerRegionName == "Екатеринбург" и отбросил бы всю
московскую выдачу целиком. offerRegionName как критерий тоже не годится:
часть офферов Новой Москвы приходит с именами вида "г. Говорово", а не
"Москва" — гард по названию региона молча резал бы легитимные лоты.
Координаты — единственный признак, который не врёт.
last_raw_count фиксируется ДО геогарда (как raw-счётчик у Циан) — иначе
страница, где все 20 items легитимно оказались вне bbox, была бы
неотличима от блока/пустого ответа в parse_page.
"""
payload = _domclick_payload(scraper, html)
scraper.last_raw_count = 0
if not payload:
return []
result = payload.get("result")
items = result.get("items") if isinstance(result, dict) else None
if not isinstance(items, list):
return []
scraper.last_raw_count = len(items)
lots = []
for item in items:
loc = item.get("location") or {}
# float() в try: прогон идёт часами без присмотра, и один оффер с мусором
# в координате уронил бы ВЕСЬ коллектор на середине коридора вместо того,
# чтобы выпасть одной карточкой. Нечисловая координата = лот без гео,
# пропускаем его так же, как отсутствующую.
try:
lat = float(loc["lat"])
lon = float(loc["lon"])
except (KeyError, TypeError, ValueError):
continue
if not (_MSK_LAT_MIN <= lat <= _MSK_LAT_MAX
and _MSK_LON_MIN <= lon <= _MSK_LON_MAX):
continue
lot = scraper._map_item(item)
if lot is not None:
lots.append(lot)
return lots
def _domclick_detect_block(html: str) -> tuple[str, str] | None:
"""QRATOR/капча DomClick: HTML вместо JSON BFF-ответа.
Порядок обязателен — сперва пробуем распарсить JSON, и только если не
вышло, ищем маркеры (тот же порядок и та же причина, что в докстринге
_extract_json в ките, #3267): маркер вроде "система защиты" или "captcha"
может встретиться в тексте самого объявления, и подстрочный поиск ДО
попытки распарсить JSON забанил живой узел на 6 часов в проде 30.08.
Валидный JSON нужной формы блок-страницей быть не может — QRATOR всегда
отдаёт HTML, поэтому успешный json.loads сам по себе доказывает «блока
нет», и маркеры уже не смотрим.
HTTP-400-подобное тело ({"statusCode":400,"error":"Bad Request"}, ответ на
offset за потолком пагинации) — валидный JSON без result.pagination, а не
блок: возвращаем None, конец коридора отловит parse_page по правилу
«count is None и 0 карточек». В штатном прогоне это тело недостижимо —
DOMCLICK_MAX_PAGES ограничивает offset числом ниже HTTP-400 порога.
"""
start = html.find("{")
end = html.rfind("}")
if start != -1 and end != -1 and end > start:
try:
json.loads(html[start:end + 1])
return None
except json.JSONDecodeError:
pass
html_lower = html.lower()
for marker in DOMCLICK_BLOCK_MARKERS:
if marker in html_lower:
return "challenge", f"маркер блока DomClick: {marker!r}"
return None
async def _domclick_wait_ready(page: Any) -> None:
"""BFF отдаёт JSON текстовым узлом в <pre> (как gate-API Яндекса) — ждём его
появления коротко; на капче его не будет вовсе, и через 3 с отработает
_domclick_detect_block.
"""
try:
await page.wait_for_selector("pre", timeout=3_000)
except Exception: # noqa: BLE001 — блок/пустая страница разбирается ниже
pass
@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] # корутина: дождаться готовности страницы
# Дефолт --target-count. Он обязан быть <= hard_cap площадки, иначе бисекция
# штатно оставляет коридоры, которые заведомо не вычитываются до конца.
default_target: int = 1500
@property
def hard_cap(self) -> int:
return self.page_size * self.max_pages
def make_scraper(self) -> Any:
if self.name == "yandex":
# У Яндекса разбор — чистые функции над gate-JSON; YandexRealtyScraper
# нужен только ради camoufox-транспорта и пула прокси, которых здесь
# нет (транспорт — вкладка Chrome владельца). Носитель состояния для
# мемоизации payload и last_raw_count — пустой namespace.
return SimpleNamespace()
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",
)
if self.name == "domclick":
# _map_item — чистая функция над dict (см. providers/domclick/serp.py),
# self._config/self._cookies там не читаются, поэтому пустой
# namespace достаточен — тот же приём, что и у CianScraper ниже.
return DomClickScraper(SimpleNamespace()) # type: ignore[arg-type]
# 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,
),
"yandex": PlatformAdapter(
name="yandex",
table="yandex_cards",
default_base_url=DEFAULT_YANDEX_BASE_URL,
page_size=YANDEX_PAGE_SIZE,
max_pages=YANDEX_MAX_PAGES,
build_url=_yandex_build_url,
extract_total_count=_yandex_extract_total,
parse_cards=_yandex_parse_cards,
detect_block=_yandex_detect_block,
wait_ready=_yandex_wait_ready,
default_target=YANDEX_TARGET_COUNT,
),
"domclick": PlatformAdapter(
name="domclick",
table="domclick_cards",
default_base_url=DEFAULT_DOMCLICK_BASE_URL,
page_size=DOMCLICK_PAGE_SIZE,
max_pages=DOMCLICK_MAX_PAGES,
build_url=_domclick_build_url,
extract_total_count=_domclick_extract_total,
parse_cards=_domclick_parse_cards,
detect_block=_domclick_detect_block,
wait_ready=_domclick_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 Error as _PwError
from playwright.async_api import async_playwright
from playwright.async_api import TimeoutError as _PwTimeout
global PlaywrightTimeoutError, PlaywrightError
PlaywrightTimeoutError = _PwTimeout
PlaywrightError = _PwError
# 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, net_attempts: int = 5):
"""goto с ограниченным ретраем на таймаут навигации.
Авито изредка держит соединение до упора и goto падает по timeout. Это
НЕ признак отказа: в наблюдавшемся случае вкладка показывала нормальную
выдачу, а маркеров фаервола/PoW не было. Но и молча ретраить бесконечно
нельзя — тихий отказ выглядит ровно так же. Поэтому: перед каждым
повтором пробуем прочитать то, что в документе, и прогнать через _guard,
чтобы настоящий блок остановил прогон с правильной причиной; исчерпали
попытки — жёсткий стоп с причиной nav_timeout.
Отдельно — сетевые отказы самой машины (`_TRANSIENT_NET_ERRORS`). Их
считаем своим счётчиком и ждём дольше: обрыв сети длится минуты, а не
секунды, и площадка тут ни при чём. Именно на таком отказе 12.09 умер
восьмичасовой проход Яндекса. Всё, что не в списке, поднимается как
раньше — молча глотать незнакомую ошибку навигации нельзя, тихий отказ
выглядит ровно так же.
"""
net_failures = 0
i = -1
while True:
i += 1
try:
return await self._page.goto(url, wait_until="domcontentloaded",
timeout=90_000)
except PlaywrightError as exc: # noqa: PERF203
if not isinstance(exc, PlaywrightTimeoutError):
if not _is_transient_net_error(exc):
raise
net_failures += 1
if net_failures >= net_attempts:
raise Blocked("nav_network") from None
wait_s = min(120.0, 15.0 * net_failures)
print(
f" сеть машины отвалилась ({str(exc).splitlines()[0][:80]}), "
f"ждём {wait_s:.0f} с, попытка {net_failures + 1}/{net_attempts}",
flush=True,
)
await asyncio.sleep(wait_s)
i -= 1 # сетевой сбой не тратит бюджет попыток по таймауту
continue
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
# offerId Яндекса — 19-значный (напр. 3937530341842304552), это уже
# впритык к signed bigint. Число сверх диапазона уронило бы \copy всего
# батча целиком, а не одну строку, поэтому отсекаем здесь же.
if not (-(2**63) <= source_id < 2**63):
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 Авито/Циан/Яндекса/DomClick (вторичка, Москва+МО) в"
" прод-схему 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=None,
help="целевой размер коридора; больше — делим "
"(дефолт зависит от --platform: 1500 у avito/cian, "
"500 у yandex — там потолок пагинации 25*20)")
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.target_count is None:
args.target_count = ADAPTERS[args.platform].default_target
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())