Some checks failed
CI Trade-In / changes (pull_request) Successful in 7s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 10s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Failing after 5m19s
Прогон 5425 оборвался по доле блоков и получил статус banned на 48 «блоках», из которых 41 был отказом нашего сайдкара: record_block() вида не принимал, поэтому infra падал в числитель скользящего окна #3184 наравне с настоящим баном, а mark_backfill_finished считал диагноз только ради телеметрии. - record_block(kind): в числитель идёт только platform; всё остальное — в знаменатель (как record_failure), мимо серии и safety-net; - NoProxyAvailableError у avito — не блок и не отказ площадки: прогон завершается no_proxy_stop=1 + mark_failed «пул прокси пуст» (как домклик после #3283); опознаётся по цепочке __cause__, не по подстроке (#3272); - статус banned — только при доминировании platform; при infra прогон получает failed (нулевой результат) или done, с честной причиной.
1098 lines
69 KiB
Python
1098 lines
69 KiB
Python
"""Scheduled backfill: detail-enrichment for legacy avito listings (#1551).
|
||
|
||
Nightly window 09:00-12:00 UTC (migration 112, source=avito_detail_backfill).
|
||
Offset from avito_city_sweep (~03-04 UTC) to avoid sharing proxy IP simultaneously.
|
||
|
||
Problem: ~15243/15917 avito listings have detail_enriched_at IS NULL.
|
||
City sweep ENRICH_DETAIL only enriches top-N proximate listings per anchor.
|
||
Legacy listings (older than 2h or outside radius) are never enriched.
|
||
|
||
Solution: single snapshot SELECT at start (guarantees termination), same proxy
|
||
session path as the detail-phase of `run_avito_city_sweep`
|
||
(scraper_kit.orchestration.pipeline). Block handling mirrors that phase:
|
||
rotate IP on every block, abort after max_consecutive_blocks. Статус оборванного
|
||
блоками прогона — 'banned' (#2674, runs.mark_backfill_finished): работу он не
|
||
доделал, остаток снапшота уедет в следующую ночь через NULL detail_enriched_at.
|
||
|
||
Отказы, не являющиеся блоками, до 2026-08-06 брейкера не имели вовсе: прогоны
|
||
3-5 августа делали ~1600 попыток, получали 1600 отказов, ноль обогащений и
|
||
выедали весь бюджет (9000 с) вместе с 1600 запросами через единственный прокси.
|
||
Теперь такая серия обрывается по max_consecutive_failures, а самая частая причина
|
||
отказа пишется в текст статуса прогона (_failure_signature) — иначе она живёт
|
||
только в логах контейнера, а те исчезают на первом же деплое.
|
||
|
||
ДИАГНОЗ бана (#2764): каждый пойманный блок классифицируется по ТИПУ исключения
|
||
(ban_kind_of_exception) и уходит в scrape_runs.ban_kind. До этой правки прогон
|
||
3306 (blocked=5, 0 обогащено) получил 'platform' по УМОЛЧАНИЮ — финализатор
|
||
диагноз не передавал, а browser-режим fetch_detail всё равно превращал отказ
|
||
сайдкара в AvitoBlockedError, так что передавать было бы нечего.
|
||
|
||
КРИТЕРИЙ ОБРЫВА (#3184): голое "N блоков подряд" (max_consecutive_blocks)
|
||
абортило прогон 125/190 (25% блоков, честно обогащённых 125 карточек) статусом
|
||
'banned' наравне с прогоном 0/5 (100% блоков) — блоки идут автокоррелированными
|
||
пачками, и пачка 5-6 подряд у здорового прогона с базовой долей 25-46% не
|
||
редкость. Теперь абортит app.services.backfill_block_breaker.BlockRatioBreaker:
|
||
доля блоков в скользящем окне последних N попыток (settings.detail_backfill_
|
||
block_ratio_window/_threshold) ИЛИ safety-net для коротких прогонов (max_
|
||
consecutive_blocks блоков подряд БЕЗ единого успеха с самого начала — старое
|
||
поведение для случая вида 0/5). Гистограмма длин пачек блоков едет в
|
||
counters["block_streak_histogram"], иначе эффект правки нечем измерить постфактум.
|
||
Какой из двух критериев сработал -- в counters["abort_reason"] ("ratio" /
|
||
"safety_net"); ключа нет, если прогон не обрывался. Он же называется в ABORT-логе
|
||
своей величиной: доля печатает "14/20", safety-net -- длину серии. Раньше лог
|
||
печатал серию всегда, и прогон 5210 (обрыв по доле 14/20) отчитался как
|
||
"ABORT -- 1 consecutive blocks".
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import asyncio
|
||
import logging
|
||
import random
|
||
import re
|
||
import time
|
||
from collections import Counter
|
||
from dataclasses import dataclass, field
|
||
from urllib.parse import urlparse
|
||
|
||
from curl_cffi.requests import AsyncSession
|
||
from scraper_kit.avito_exceptions import (
|
||
AvitoBlockedError,
|
||
AvitoListingGoneError,
|
||
AvitoRateLimitedError,
|
||
)
|
||
from scraper_kit.browser_fetcher import BrowserFetcher
|
||
from scraper_kit.orchestration.pipeline import CITY_LOCATIONS, ban_kind_of_exception
|
||
|
||
# #2397 slice B (эпик #2277 decommission scrape_pipeline.py, Part E): раньше
|
||
# _CHROME_HEADERS/_avito_proxies() импортировались из app.services.scrape_pipeline.
|
||
# kit уже вынес byte-идентичный dict (DOCUMENT_HEADERS) и ту же формулу
|
||
# {"http": url, "https": url} if url else None (http_proxies) в
|
||
# providers/_base.py (#2358 Foundation) — реюзаем вместо копии, разрывая
|
||
# зависимость от scrape_pipeline.py перед его удалением.
|
||
from scraper_kit.providers._base import DEFAULT_IMPERSONATE, DOCUMENT_HEADERS, http_proxies
|
||
from scraper_kit.providers.avito.detail import (
|
||
_AVITO_WARM_SEARCH_URL,
|
||
_serp_origin_for,
|
||
build_warmed_session,
|
||
fetch_detail,
|
||
research_in_session,
|
||
save_detail_enrichment,
|
||
)
|
||
from scraper_kit.providers.avito.serp import AvitoScraper
|
||
from scraper_kit.proxy_errors import NoProxyAvailableError
|
||
from scraper_kit.snapshot_writer import upsert_listing_snapshot
|
||
from sqlalchemy import text
|
||
from sqlalchemy.orm import Session
|
||
|
||
from app.core.config import settings
|
||
from app.core.shutdown import shutdown_requested
|
||
from app.services import scrape_runs as runs_mod
|
||
from app.services.backfill_block_breaker import BlockRatioBreaker
|
||
from app.services.proxy_egress import resolve_proxy_url
|
||
from app.services.proxy_rotation import rotate_proxy
|
||
from app.services.scraper_adapters import RealProxyProvider, RealScraperConfig
|
||
|
||
# #2397 Part D1 (#2330 закрыт): _AVITO_WARM_SEARCH_URL/build_warmed_session больше
|
||
# НЕ legacy — kit's scraper_kit.providers.avito.detail.build_warmed_session() теперь
|
||
# принимает config: ScraperConfig | None (Strangler-инъекция #2330) и пробрасывает
|
||
# его в _build_detail_session() для settings.scraper_proxy_url (sticky МГТС-прокси),
|
||
# зеркаля fetch_detail. Оба call site'а ниже передают config=RealScraperConfig() —
|
||
# без него warm-batch curl-путь (use_curl=True, ДЕФОЛТ прод-режима,
|
||
# avito_detail_backfill_use_curl=True) молча терял бы backconnect (прямое
|
||
# datacenter-подключение вместо sticky-прокси), тот же класс бага что #2322/#2310.
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
__all__ = [
|
||
"AvitoDetailBackfillResult",
|
||
"run_avito_detail_backfill",
|
||
]
|
||
|
||
# #2576 этап B: oblast-города (region 66, вне ЕКБ) уже дают листинги (Каменск-
|
||
# Уральский), но snapshot-SELECT ниже раньше фильтровал ЖЁСТКО '%/ekaterinburg/%' —
|
||
# у всех остальных detail_enriched_at оставался NULL навсегда (без detail-страницы
|
||
# нет lat/lon -> листинг молча выпадает из подбора аналогов по радиусу).
|
||
# CITY_LOCATIONS.avito_slug — единственный источник правды для avito URL-слага
|
||
# города (может отличаться от нашего city_slug: kamensk-uralskiy через дефис,
|
||
# verhnyaya_pyshma без "kh") -- дублировать список тут вместо импорта было бы
|
||
# risk дрейфа при добавлении новых oblast-городов.
|
||
#
|
||
# #2578 review: Postgres LIKE трактует '_' как wildcard "один любой символ" (не
|
||
# литерал) и '%' как wildcard "любая последовательность" -- два слага из пяти
|
||
# (nizhniy_tagil, verhnyaya_pyshma) содержат '_', без экранирования это латентная
|
||
# дыра: город с похожим слагом (напр. nizhniyXtagil) молча совпал бы. Сегодня
|
||
# коллизий нет (проверено на проде: raw vs escaped паттерны дают одинаковые 776
|
||
# совпадений), но экранируем сейчас, а не когда появится реальная коллизия.
|
||
# LIKE по умолчанию использует '\' как escape-символ (без явного ESCAPE) —
|
||
# подтверждено на живом Postgres 16.4 (см. коммит #2578-fixup): 'nizhniyXtagil'
|
||
# матчит неэкранированный '%/nizhniy_tagil/%' (LIKE default '_'=wildcard) и НЕ
|
||
# матчит экранированный '%/nizhniy\_tagil/%' (LIKE '\_' = литерал '_'); точный
|
||
# слаг 'nizhniy_tagil' матчит оба варианта -- позитивный кейс не сломан.
|
||
#
|
||
# #262 wave 2: avito_slug — Optional в CityLocation (не у каждого oblast-города
|
||
# подтверждён). Города без avito_slug пропускаем целиком — у них НЕТ avito_city_
|
||
# sweep schedule (262_), значит НЕТ и avito-листингов с их URL; паттерн для них
|
||
# был бы либо мёртвым, либо (что хуже) построен из city_slug вместо реального
|
||
# avito URL-сегмента и создал бы ложный LIKE-матч.
|
||
_OBLAST_AVITO_URL_PATTERNS = tuple(
|
||
"%/" + loc.avito_slug.replace("\\", "\\\\").replace("_", "\\_").replace("%", "\\%") + "/%"
|
||
for loc in CITY_LOCATIONS.values()
|
||
if loc.avito_slug is not None
|
||
)
|
||
|
||
|
||
# Причина отказа карточки без её URL: 1576 отказов одного прогона должны схлопнуться
|
||
# в ОДНУ строку, иначе перепись бесполезна.
|
||
_URL_IN_MESSAGE_RE = re.compile(r"https?://\S+")
|
||
|
||
|
||
def _failure_signature(exc: BaseException) -> str:
|
||
"""Подпись причины отказа: тип исключения + текст без URL.
|
||
|
||
Зачем (замер 2026-08-06): у прогонов 3 и 4 августа counters говорили
|
||
`attempted=1576, failed=1576, blocked=0` — и ничего больше. Кто отказал,
|
||
площадка или наш тракт, было видно ТОЛЬКО в логах контейнера, а тот
|
||
пересоздаётся на каждом деплое и уносит их с собой; в GlitchTip попадают
|
||
события уровня ERROR, а поштучные отказы — WARNING. Разница между этими
|
||
двумя диагнозами — разные владельцы задачи (#2686, #2698), поэтому она
|
||
обязана переживать перезапуск контейнера, то есть лежать в самом прогоне.
|
||
|
||
Тип исключения — первый разряд диагноза (AvitoBlockedError = площадка
|
||
показала 403/firewall; сетевой класс curl_cffi = наш прокси-тракт;
|
||
ValueError = ответ пришёл, но не разобран), текст — второй.
|
||
"""
|
||
message = _URL_IN_MESSAGE_RE.sub("<url>", str(exc)).strip()
|
||
return f"{type(exc).__name__}: {message}"[:160] if message else type(exc).__name__
|
||
|
||
|
||
def _top_failure(census: Counter[str]) -> str | None:
|
||
"""Самая частая причина отказа с её долей; None — отказов не было."""
|
||
if not census:
|
||
return None
|
||
reason, hits = census.most_common(1)[0]
|
||
return f"{reason} ({hits} из {sum(census.values())})"
|
||
|
||
|
||
def _iter_causes(exc: BaseException) -> list[BaseException]:
|
||
"""Цепочка причин исключения, без зацикливания (копия приёма из
|
||
domclick_detail_backfill — там же он и обкатан на #3283)."""
|
||
seen: set[int] = set()
|
||
out: list[BaseException] = []
|
||
cur: BaseException | None = exc
|
||
while cur is not None and id(cur) not in seen:
|
||
out.append(cur)
|
||
seen.add(id(cur))
|
||
cur = cur.__cause__ or cur.__context__
|
||
return out
|
||
|
||
|
||
def _caused_by_empty_pool(exc: BaseException) -> bool:
|
||
"""Прячется ли за этим «блоком» пустой пул прокси (#3288, как #3283 у домклика).
|
||
|
||
`NoProxyAvailableError` документирован ровно как «НАША инфраструктура, не
|
||
внешний блок», и поднимается ДО HTTP-запроса: к площадке мы не ходили вовсе.
|
||
Сюда он попадает под видом блокировки, потому что fetch_detail заворачивает
|
||
в `AvitoSidecarUnavailableError` любое исключение фетча.
|
||
|
||
Опора — ТИП в цепочке `__cause__`/`__context__`, а не подстрока «no proxy
|
||
available» в тексте: текст обёртки её действительно содержит, но ровно так
|
||
#3272 уже один раз объявил блоком пользовательское описание квартиры.
|
||
"""
|
||
return any(isinstance(c, NoProxyAvailableError) for c in _iter_causes(exc))
|
||
|
||
|
||
# Провайдер mobileproxy.space почти всегда возвращает "rt" (секунды на переподключение
|
||
# канала после ротации) в RotationResult.reconnect_delay_s -- см. app.services.
|
||
# proxy_rotation docstring. Дефолт нужен ТОЛЬКО если провайдер его не прислал
|
||
# (best-effort парсинг "rt", неудача не является ошибкой ротации) -- замеренный
|
||
# вживую диапазон простоя канала 2-12с, 10с чуть выше нижней границы.
|
||
_DEFAULT_ROTATE_RECONNECT_DELAY_S = 10.0
|
||
|
||
|
||
async def _rotate_current_proxy(
|
||
db: Session,
|
||
run_id: int,
|
||
browser_fetcher: BrowserFetcher,
|
||
trigger: str = "rotate-by-attempts",
|
||
) -> bool:
|
||
"""Ротация exit-IP арендованного прокси + сброс browser-контекста.
|
||
|
||
Два триггера, оба пишут свой ярлык в логи через ``trigger``: плановый по счётчику
|
||
попыток (``rotate-by-attempts``) и реактивный по бану площадки (``rotate-on-ban``,
|
||
#3283). Разделение нужно только для читаемости логов — механика одна и та же.
|
||
|
||
Свежий адрес со старыми куками бесполезен — личность браузера должна меняться
|
||
ВМЕСТЕ с адресом, иначе следующий запрос уходит с нового IP, но с cookie-следом
|
||
старого (request_context_reset потребляется ровно следующим fetch(), #3118).
|
||
|
||
Работает ТОЛЬКО когда есть browser_fetcher (browser-режим) — это единственный
|
||
путь, где avito_detail_backfill реально держит lease с proxy_id (scrape_proxies.id,
|
||
browser_fetcher.lease_id). curl/backconnect-путь (use_curl=True) идёт через sticky
|
||
settings.scraper_proxy_url БЕЗ lease — там ротировать нечего, вызывающий цикл не
|
||
зовёт эту функцию в том режиме вовсе.
|
||
|
||
Прод при этом в browser-режиме, а не в curl: у контейнера tradein-scraper (там же
|
||
живёт планировщик) проверено AVITO_DETAIL_BACKFILL_USE_CURL=false при
|
||
SCRAPER_FETCH_MODE=browser и USE_PROXY_POOL_BROWSER=true, так что ротация
|
||
активна. Значение true стоит только у tradein-backend, который добор не запускает.
|
||
Комментарий ниже по файлу (~строка 653) называет use_curl=True «прод-дефолтом» —
|
||
это предсуществующее заблуждение, а не описание текущего прода.
|
||
|
||
Отказ провайдера (лимит исчерпан, нет rotate_url, сетевой сбой, неизвестный хост)
|
||
НЕ должен ронять прогон — логируем и продолжаем на текущем адресе. Эта попытка НЕ
|
||
является блоком/отказом площадки и не должна попадать в BlockRatioBreaker/
|
||
counters.blocked/ban_kinds — вызывающий код не передаёт её исход ни в один из них,
|
||
это обслуживание канала, а не результат fetch().
|
||
"""
|
||
proxy_id = browser_fetcher.lease_id
|
||
if proxy_id is None:
|
||
logger.info(
|
||
"avito_detail_backfill: run_id=%d %s skipped -- no leased proxy",
|
||
run_id,
|
||
trigger,
|
||
)
|
||
return False
|
||
|
||
result = await rotate_proxy(db, proxy_id)
|
||
if not result.ok:
|
||
logger.warning(
|
||
"avito_detail_backfill: run_id=%d %s FAILED proxy_id=%d: %s",
|
||
run_id,
|
||
trigger,
|
||
proxy_id,
|
||
result.reason,
|
||
)
|
||
return False
|
||
|
||
delay = (
|
||
result.reconnect_delay_s
|
||
if result.reconnect_delay_s is not None
|
||
else _DEFAULT_ROTATE_RECONNECT_DELAY_S
|
||
)
|
||
logger.info(
|
||
"avito_detail_backfill: run_id=%d %s OK proxy_id=%d new_ip=%s -- "
|
||
"waiting %.1fs for channel reconnect",
|
||
run_id,
|
||
trigger,
|
||
proxy_id,
|
||
result.new_ip,
|
||
delay,
|
||
)
|
||
await asyncio.sleep(delay)
|
||
browser_fetcher.request_context_reset()
|
||
return True
|
||
|
||
|
||
@dataclass
|
||
class AvitoDetailBackfillResult:
|
||
"""Counters for one backfill run."""
|
||
|
||
attempted: int = 0
|
||
enriched: int = 0
|
||
blocked: int = 0
|
||
gone: int = 0
|
||
failed: int = 0
|
||
duration_sec: float = field(default=0.0)
|
||
|
||
def to_dict(self) -> dict[str, int]:
|
||
return {
|
||
"attempted": self.attempted,
|
||
"enriched": self.enriched,
|
||
"blocked": self.blocked,
|
||
"gone": self.gone,
|
||
"failed": self.failed,
|
||
"duration_sec": int(self.duration_sec),
|
||
}
|
||
|
||
|
||
async def run_avito_detail_backfill(
|
||
db: Session,
|
||
*,
|
||
run_id: int,
|
||
params: dict,
|
||
) -> AvitoDetailBackfillResult:
|
||
"""Backfill detail_enriched_at for legacy avito listings via mobile proxy.
|
||
|
||
Params (from default_params jsonb in scrape_schedules):
|
||
batch_size: int -- ЕКБ snapshot size (SELECT LIMIT), default 800 (unchanged,
|
||
#2576 -- volume/order for ЕКБ stay byte-identical to pre-oblast behaviour).
|
||
oblast_batch_size: int -- ДОПОЛНИТЕЛЬНАЯ reserved-квота для листингов
|
||
области (#2576), default 100. Отдельный LIMIT, НЕ отъедает от batch_size
|
||
ЕКБ -- гарантирует области честную обработку и одновременно не даёт
|
||
всплеску свежих oblast-листингов вытеснить ЕКБ из top-N по scraped_at.
|
||
budget_sec: float -- wall-clock budget per run, default 3600s.
|
||
request_delay_sec: float -- delay between listings, default 6.0s.
|
||
max_consecutive_blocks: int -- safety_min для BlockRatioBreaker (#3184):
|
||
абортит прогоны короче окна, если ВСЕ попытки с начала были блоками
|
||
(ни одного успеха), default 5. Основной критерий -- доля блоков в
|
||
скользящем окне, см. settings.detail_backfill_block_ratio_window/
|
||
_threshold (module docstring).
|
||
max_consecutive_failures: int -- порог обрыва по отказам-не-блокам,
|
||
default 25 (см. комментарий у чтения параметра ниже).
|
||
|
||
Lifecycle: update_heartbeat -> snapshot -> loop with budget guard ->
|
||
mark_backfill_finished (done / banned при блоках / failed при нуле, #2674);
|
||
mark_failed напрямую — только при исключении.
|
||
"""
|
||
batch_size = int(params.get("batch_size", 800))
|
||
oblast_batch_size = int(params.get("oblast_batch_size", 100))
|
||
budget_sec = float(params.get("budget_sec", 3600))
|
||
request_delay_sec = float(params.get("request_delay_sec", 6.0))
|
||
# #3184: теперь safety_min BlockRatioBreaker (см. module docstring) -- не общий
|
||
# критерий обрыва, а страховка для прогонов короче окна.
|
||
max_consecutive_blocks = int(params.get("max_consecutive_blocks", 5))
|
||
# Брейкер на отказы-НЕ-блоки. Блоки свой брейкер имели с самого начала, отказы —
|
||
# нет, и это стоило трёх ночей подряд: 3-5 августа прогон делал ~1600 попыток,
|
||
# получал 1600 отказов, ноль обогащений и выедал весь бюджет 9000 с (плюс 1600
|
||
# запросов через единственный прокси, #2638). Порог заметно выше блочного: пачка
|
||
# мёртвых карточек (404 → ValueError в curl-режиме) не должна обрывать здоровый
|
||
# прогон, а 25 отказов подряд без единого успеха — уже не невезение.
|
||
max_consecutive_failures = int(params.get("max_consecutive_failures", 25))
|
||
warm_batch = int(params.get("warm_batch", 500))
|
||
research_every = int(params.get("research_every", 50))
|
||
block_cooldown_sec = float(params.get("block_cooldown_sec", 30.0))
|
||
warm_timeout_s = float(params.get("warm_timeout_s", 90.0))
|
||
|
||
# #1950: hard-timeout'ы на блокирующие await'ы внутри loop'а. Без них один
|
||
# зависший fetch_detail/rotate ронял весь run в zombie (run 423 завис 7.7ч).
|
||
# Локальные float'ы (а не settings.* напрямую) — чтобы wait_for/format получали
|
||
# число даже когда settings замокан MagicMock'ом в тестах.
|
||
fetch_timeout_s = float(settings.avito_detail_fetch_timeout_s)
|
||
rotate_timeout_s = float(settings.avito_proxy_rotate_settle_s) + 30.0
|
||
|
||
counters = AvitoDetailBackfillResult()
|
||
current_counters: dict[str, int] = counters.to_dict()
|
||
|
||
# Режим фетча: curl_cffi+backconnect (mproxy) когда use_curl=True; иначе браузер/auv.
|
||
# avito_detail_backfill_use_curl=True переопределяет scraper_fetch_mode — detail_backfill
|
||
# всегда идёт через backconnect чтобы не конкурировать за auv с SERP (city_sweep/full_load).
|
||
use_curl = settings.avito_detail_backfill_use_curl
|
||
browser_mode = settings.scraper_fetch_mode == "browser" and not use_curl
|
||
session: AsyncSession | None = None
|
||
browser_fetcher: BrowserFetcher | None = None
|
||
own_session = False
|
||
own_browser = False
|
||
|
||
# kit AvitoScraper требует ScraperConfig позиционно (Strangler-инжекция #2133) —
|
||
# RealScraperConfig проксирует settings.* так же, как читал legacy-конструктор без
|
||
# аргументов (scraper_proxy_url и т.д. для _build_cffi_session()/_rotate_ip()).
|
||
scraper = AvitoScraper(RealScraperConfig())
|
||
start = time.monotonic()
|
||
|
||
logger.info(
|
||
"avito_detail_backfill: run_id=%d mode=%s",
|
||
run_id,
|
||
"curl/backconnect" if use_curl else "browser",
|
||
)
|
||
|
||
try:
|
||
# Setup session (mirrors run_avito_pipeline lines 161-183).
|
||
# use_curl=True: session=None + browser_fetcher=None → fetch_detail строит
|
||
# эфемерную _build_detail_session() через settings.scraper_proxy_url (backconnect).
|
||
# use_curl=False (legacy): browser_mode=True → BrowserFetcher/auv, как раньше.
|
||
if browser_mode:
|
||
# proxy_provider/use_pool/environment — обязательная часть проводки,
|
||
# а не опция (#2698). Без них BrowserFetcher не кладёт "proxy" в тело
|
||
# POST /fetch, сайдкар берёт свой env-прокси, и прогон уходит мимо
|
||
# пула: ни выбора узла по affinity, ни учёта scrape_proxy_source_bans,
|
||
# ни ротации при блоке. В логе это видно как proxy_lease_id=None.
|
||
#
|
||
# Ровно этот дефект уже чинили в house_imv_backfill (#2698) — там он
|
||
# держал 35 из 35 отказов в каждом прогоне полтора месяца, пока
|
||
# соседние свипы через ТОТ ЖЕ сайдкар тянули сотни объявлений. Здесь
|
||
# он остался.
|
||
#
|
||
# Замер 27.08: прогон 5098 — mode=browser, proxy_lease_id=None,
|
||
# 5 блоков подряд из 5 попыток, enriched=0. При этом пул здоров
|
||
# (4 узла, все ok), а тот же URL через прокси отдаёт 200 и 3.3 МБ.
|
||
# reuse_context=True (#3180/#3251): без него sidecar's browser.new_page()
|
||
# создаёт НОВЫЙ изолированный context на КАЖДЫЙ /fetch — пройденный
|
||
# QRATOR-PoW предыдущей карточки выбрасывается, и следующий запрос снова
|
||
# холодный. У DomClick (#3118) это измеренно давало 100% блоков (26/26)
|
||
# на изолированных контекстах против 5/5 успехов в тёплом контексте.
|
||
# Сброс сожжённого context'а — ниже, через bf.request_context_reset()
|
||
# (#3212: один раз за прогон, не на каждый блок — см. except-ветку).
|
||
_cfg = RealScraperConfig()
|
||
browser_fetcher = BrowserFetcher(
|
||
source="avito",
|
||
endpoint=settings.browser_http_endpoint,
|
||
proxy_provider=RealProxyProvider(),
|
||
use_pool=_cfg.use_proxy_pool_browser,
|
||
# Без environment отказ «пул пуст» на этом пути мёртв — фетчер
|
||
# молча ушёл бы на env-прокси сайдкара (#2616 шаг 1).
|
||
environment=_cfg.environment,
|
||
reuse_context=True,
|
||
)
|
||
await browser_fetcher.__aenter__()
|
||
own_browser = True
|
||
scraper._browser = browser_fetcher
|
||
elif use_curl:
|
||
# use_curl=True (#1551 warm-batch): прогретая shared-сессия на sticky МГТС-IP —
|
||
# yandex-referer -> avito-search сеет антибот-куки, сессия держит батч detail без
|
||
# 403 (доказано пробами: ≥66 карточек подряд, 0 блоков). NB: МГТС sticky — один
|
||
# фикс. exit-IP, per-connection ротации нет (rebuild != новый IP); on-block —
|
||
# cooldown + in-session re-search, не дискард сессии.
|
||
own_session = True
|
||
# config= обязателен: #2397 Part D1 — kit build_warmed_session() без него
|
||
# молча теряет settings.scraper_proxy_url (sticky МГТС-прокси), см. модуль-
|
||
# ный комментарий у импорта выше.
|
||
session = await build_warmed_session(config=RealScraperConfig())
|
||
elif not use_curl:
|
||
# curl_cffi legacy path (scraper_fetch_mode="curl_cffi", use_curl=False):
|
||
# строим shared сессию через auv, как делает run_avito_city_sweep (kit).
|
||
# Резолвер по источнику (#2825): пул scrape_proxies с учётом
|
||
# scrape_proxy_source_bans, fallback на settings.scraper_proxy_url только
|
||
# если пул пуст (легитимный dev/staging-сценарий). Пул не пуст, но все
|
||
# забанены/нездоровы для avito -- resolve_proxy_url бросает
|
||
# ProxyPoolExhaustedError (fail-closed, #2616): НАРОЧНО не ловим здесь --
|
||
# штатный except Exception ниже (mark_failed + logger.exception + raise)
|
||
# уже даёт явную деградацию run'а с понятным логом, отдельный catch не нужен.
|
||
own_session = True
|
||
session = AsyncSession(
|
||
impersonate=DEFAULT_IMPERSONATE,
|
||
timeout=25,
|
||
headers=DOCUMENT_HEADERS,
|
||
proxies=http_proxies(resolve_proxy_url(db, "avito")),
|
||
)
|
||
scraper._cffi = session
|
||
|
||
runs_mod.update_heartbeat(db, run_id, current_counters)
|
||
|
||
# SNAPSHOT: single SELECT at start -- NOT re-selected in loop.
|
||
# Scope (#1814, расширено #2576): активные листинги ЕКБ + известных oblast-
|
||
# городов (region 66). region_code на insert хардкодится в 66 (base.py) →
|
||
# НЕ дискриминирует город; реальный признак города у Avito — путь URL
|
||
# (/ekaterinburg/ для ЕКБ; legacy Москва/СПб/Тюмень — /moskva//sankt-
|
||
# peterburg//tyumen/ — те по-прежнему вне scope, НЕ входят ни в ekb, ни в
|
||
# oblast CTE). browser-fetch на legacy не-ЕКБ/не-oblast спотыкается →
|
||
# curl-fallback → 429-бан curl-фингерпринта. Не тратим фетчи на мёртвые
|
||
# (is_active) и на регионы вне scope.
|
||
#
|
||
# Два CTE вместо одного WHERE ... OR ...: ekb сохраняет ТОЧНО прежний
|
||
# LIMIT/ORDER (#2576 требование "ЕКБ не деградирует") -- oblast НЕ может
|
||
# вытеснить ЕКБ из batch_size ни при каком всплеске свежих oblast-строк
|
||
# (ORDER BY ... scraped_at DESC в общем WHERE отдал бы приоритет самым
|
||
# свежим независимо от города). oblast получает отдельную честную квоту
|
||
# oblast_batch_size, добавленную ПОСЛЕ ekb-квоты (не вычтенную из неё).
|
||
snapshot = (
|
||
db.execute(
|
||
text(
|
||
"""
|
||
WITH ekb AS (
|
||
SELECT id, source_url, price_rub, 'ekb' AS city_scope
|
||
FROM listings
|
||
WHERE source = 'avito'
|
||
AND detail_enriched_at IS NULL
|
||
AND source_url IS NOT NULL
|
||
AND is_active = TRUE
|
||
AND source_url LIKE '%/ekaterinburg/%'
|
||
-- сперва листинги без координат (#1967 — detail-страница
|
||
-- даёт координаты здания), затем по свежести
|
||
ORDER BY (lat IS NULL) DESC, scraped_at DESC NULLS LAST
|
||
LIMIT CAST(:batch_size AS int)
|
||
),
|
||
oblast AS (
|
||
SELECT id, source_url, price_rub, 'oblast' AS city_scope
|
||
FROM listings
|
||
WHERE source = 'avito'
|
||
AND detail_enriched_at IS NULL
|
||
AND source_url IS NOT NULL
|
||
AND is_active = TRUE
|
||
AND source_url LIKE ANY(CAST(:oblast_patterns AS text[]))
|
||
ORDER BY (lat IS NULL) DESC, scraped_at DESC NULLS LAST
|
||
LIMIT CAST(:oblast_batch_size AS int)
|
||
)
|
||
SELECT id, source_url, price_rub, city_scope FROM ekb
|
||
UNION ALL
|
||
SELECT id, source_url, price_rub, city_scope FROM oblast
|
||
"""
|
||
),
|
||
{
|
||
"batch_size": batch_size,
|
||
"oblast_patterns": list(_OBLAST_AVITO_URL_PATTERNS),
|
||
"oblast_batch_size": oblast_batch_size,
|
||
},
|
||
)
|
||
.mappings()
|
||
.all()
|
||
)
|
||
|
||
if not snapshot:
|
||
logger.info(
|
||
"avito_detail_backfill: run_id=%d -- no pending listings "
|
||
"(detail_enriched_at IS NULL = 0), done",
|
||
run_id,
|
||
)
|
||
runs_mod.mark_done(db, run_id, current_counters)
|
||
return counters
|
||
|
||
# #2576: разбивка ekb/oblast только для наблюдаемости -- .get() консервативен
|
||
# (city_scope нет в mock-снапшотах старых тестов, дефолт "ekb" их не ломает).
|
||
oblast_count = sum(1 for row in snapshot if row.get("city_scope") == "oblast")
|
||
logger.info(
|
||
"avito_detail_backfill: run_id=%d snapshot=%d (ekb=%d oblast=%d, "
|
||
"budget=%.0fs delay=%.1fs max_blocks=%d mode=%s)",
|
||
run_id,
|
||
len(snapshot),
|
||
len(snapshot) - oblast_count,
|
||
oblast_count,
|
||
budget_sec,
|
||
request_delay_sec,
|
||
max_consecutive_blocks,
|
||
"curl/backconnect" if use_curl else "browser",
|
||
)
|
||
|
||
# #3184: доля блоков в скользящем окне вместо голого "N подряд" -- см. module
|
||
# docstring и app.services.backfill_block_breaker.
|
||
breaker = BlockRatioBreaker(
|
||
window_size=int(settings.detail_backfill_block_ratio_window),
|
||
ratio_threshold=float(settings.detail_backfill_block_ratio_threshold),
|
||
safety_min=max_consecutive_blocks,
|
||
# snapshot_size гейтит safety-net (#3184 review MAJOR 2): пачка блоков в
|
||
# начале ДЛИННОГО прогона не должна абортить его так же, как раньше --
|
||
# safety-net включён только когда снапшот короче окна и ratio-критерий
|
||
# физически недостижим (см. backfill_block_breaker module docstring).
|
||
snapshot_size=len(snapshot),
|
||
)
|
||
consecutive_failures = 0
|
||
aborted_by_blocks = False
|
||
abort_reason: str | None = None
|
||
# #3288: прогон оборван пустым пулом прокси — «нечем ходить», а не бан.
|
||
no_proxy_stop = False
|
||
do_sleep = False
|
||
items_since_warm = 0
|
||
# Счётчик попыток (успех ИЛИ отказ — оба тратят бюджет IP одинаково, см.
|
||
# settings.avito_detail_backfill_rotate_after_attempts) на ТЕКУЩЕМ арендованном
|
||
# прокси. Обнуляется на каждой ротации (успешной ИЛИ неуспешной — иначе
|
||
# исчерпанный дневной лимит провайдера дёргал бы rotate_proxy на каждой попытке).
|
||
attempts_since_rotation = 0
|
||
# #3251: сброс переиспользуемого browser-context'а (reuse_context=True выше)
|
||
# разрешён РОВНО один раз за прогон — зеркалит domclick_detail_backfill (#3212).
|
||
# Сброс на КАЖДЫЙ блок сам себя поддерживает: пройденный QRATOR-PoW живёт в
|
||
# context'е, сброс его выбрасывает, повторная проверка с того же IP сразу
|
||
# после принятой снова блокируется — одна осечка превращается в необратимый
|
||
# каскад блоков (см. except-ветку ниже).
|
||
context_reset_used = False
|
||
# #3283g: сколько раз ЗА ПРОГОН уже сработала ротация-на-бан (в отличие от
|
||
# context_reset_used эта ротация меняет IP вместе со сбросом, поэтому не
|
||
# подвержена каскаду #3251 и может срабатывать несколько раз, но бюджетно).
|
||
# Gap между любыми ротациями (по счётчику ИЛИ по бану) считается через
|
||
# attempts_since_rotation — он и так обнуляется на каждой ротации.
|
||
rotate_on_ban_used = 0
|
||
# Перепись причин (блоки + отказы) — переживает пересоздание контейнера,
|
||
# в отличие от логов; см. _failure_signature.
|
||
failure_census: Counter[str] = Counter()
|
||
# #2764/#3178: диагнозы всех блоков прогона по ТИПУ исключения, С кратностями
|
||
# (Counter, не set) — set терял их до решения: 5 прогонов подряд 4×platform+
|
||
# 1×infra и настоящий 2+2 приходили в mark_backfill_finished одинаково и
|
||
# получали 'unknown' оба раза, хотя первый явно платформенный. Перепись
|
||
# (kind -> count) решает _dominant_ban_kind в scrape_runs.py.
|
||
block_ban_kinds: Counter[str] = Counter()
|
||
|
||
for idx, row in enumerate(snapshot):
|
||
# Budget guard
|
||
elapsed = time.monotonic() - start
|
||
if elapsed > budget_sec:
|
||
logger.info(
|
||
"avito_detail_backfill: run_id=%d -- budget %.0fs exhausted "
|
||
"(elapsed=%.1fs), stopping at #%d/%d",
|
||
run_id,
|
||
budget_sec,
|
||
elapsed,
|
||
idx,
|
||
len(snapshot),
|
||
)
|
||
break
|
||
|
||
# #1182 Phase 2: кооперативный SIGTERM-drain (деплой recreate scraper).
|
||
# Чисто выходим на границе карточки — каждая карточка уже закоммичена,
|
||
# mark_done с текущими счётчиками вызовется после loop'а. Snapshot —
|
||
# pending-query (detail_enriched_at IS NULL), per-item идемпотентна →
|
||
# следующий run сам до-резюмит остаток, resume_cursor не нужен.
|
||
if shutdown_requested():
|
||
logger.info(
|
||
"avito_detail_backfill: run_id=%d SIGTERM-drain — stopping at #%d/%d",
|
||
run_id,
|
||
idx,
|
||
len(snapshot),
|
||
)
|
||
break
|
||
|
||
# Delay before each request except the first (джиттер ±30% — менее
|
||
# роботизированный паттерн на органически прогретой сессии).
|
||
if do_sleep:
|
||
await asyncio.sleep(request_delay_sec * random.uniform(0.7, 1.4))
|
||
do_sleep = True
|
||
|
||
source_url: str = row["source_url"]
|
||
counters.attempted += 1
|
||
attempts_since_rotation += 1
|
||
|
||
# Normalise URL -> path (mirrors run_avito_city_sweep detail-phase, kit)
|
||
item_url = urlparse(source_url).path if source_url.startswith("http") else source_url
|
||
|
||
# use_curl warm-batch (#1551): интервалы по размеру SERP-страницы (~50).
|
||
# МГТС sticky-IP: exit-IP константа (rebuild НЕ даёт новый IP). Каждые
|
||
# warm_batch (~500) — полный re-warm (close+build) на ТОМ ЖЕ sticky exit-IP:
|
||
# session-hygiene (свежий TLS-хендшейк + свежие куки yandex->search), НЕ смена
|
||
# IP. Между ними — лёгкий in-session перепоиск каждые research_every (~50)
|
||
# (re-GET search той же сессией, освежить куки). On-block — cooldown ниже.
|
||
if use_curl and session is not None:
|
||
if warm_batch and items_since_warm >= warm_batch:
|
||
# периодический полный re-warm на том же sticky exit-IP (свежий TLS+куки).
|
||
# #1950: bound сетевой re-warm; строим НОВУЮ сессию до close старой — на
|
||
# timeout/fail сохраняем текущую рабочую сессию (graceful degrade).
|
||
try:
|
||
# config= обязателен здесь тоже (тот же footgun, что и в setup-
|
||
# вызове build_warmed_session выше) -- иначе периодический re-warm
|
||
# молча пересобрал бы сессию БЕЗ sticky-прокси после первого раза.
|
||
new_session = await asyncio.wait_for(
|
||
build_warmed_session(config=RealScraperConfig()),
|
||
timeout=warm_timeout_s,
|
||
)
|
||
try:
|
||
await session.close()
|
||
except Exception:
|
||
pass
|
||
session = new_session
|
||
except Exception:
|
||
logger.warning(
|
||
"avito_detail_backfill: run_id=%d re-warm timed out/failed (>%.0fs) "
|
||
"-- keeping current session",
|
||
run_id,
|
||
warm_timeout_s,
|
||
exc_info=True,
|
||
)
|
||
items_since_warm = 0
|
||
elif (
|
||
research_every
|
||
and items_since_warm > 0
|
||
and items_since_warm % research_every == 0
|
||
):
|
||
# лёгкий in-session перепоиск: тот же exit-IP, освежить куки.
|
||
# #1950: bound сетевой re-search.
|
||
try:
|
||
await asyncio.wait_for(research_in_session(session), timeout=warm_timeout_s)
|
||
except Exception:
|
||
logger.warning(
|
||
"avito_detail_backfill: run_id=%d in-session re-search "
|
||
"timed out/failed",
|
||
run_id,
|
||
exc_info=True,
|
||
)
|
||
|
||
try:
|
||
# #1950: hard-timeout — зависший fetch (browser hang / curl-stall) не
|
||
# должен блокировать loop навсегда (иначе budget-guard/heartbeat молчат
|
||
# и run reaped как zombie). wait_for отменяет fetch → TimeoutError.
|
||
# #3251: same-site SERP-якорь. Смысл он имеет ТОЛЬКО в browser-
|
||
# режиме: сайдкар держит выдачу открытой якорной вкладкой и шлёт
|
||
# её как Referer целевой навигации. На curl/own-session пути
|
||
# сессия прогрета warm-batch'ем (см. referer= ниже, это ДРУГОЕ
|
||
# поле и другой фетчер) — там origin остаётся None и fetch_detail
|
||
# ведёт себя ровно как раньше. None также если source_url не
|
||
# распарсился на город+категорию (см. docstring _serp_origin_for).
|
||
serp_origin = _serp_origin_for(source_url) if browser_mode else None
|
||
enrichment = await asyncio.wait_for(
|
||
fetch_detail(
|
||
item_url,
|
||
cffi_session=session,
|
||
browser_fetcher=browser_fetcher,
|
||
referer=_AVITO_WARM_SEARCH_URL if use_curl else None,
|
||
reconnect_on_block=not use_curl,
|
||
# config= обязателен: kit fetch_detail без него читает
|
||
# config.scraper_proxy_url=None для backconnect-gate'а
|
||
# (elif not use_curl own-session path), теряя reconnect-
|
||
# on-403 поведение legacy (settings.scraper_proxy_url
|
||
# напрямую). Не влияет на use_curl=True (прод-дефолт) —
|
||
# там reconnect_on_block=False уже гасит backconnect.
|
||
config=RealScraperConfig(),
|
||
origin=serp_origin,
|
||
browser_referer=serp_origin,
|
||
),
|
||
timeout=fetch_timeout_s,
|
||
)
|
||
if save_detail_enrichment(db, enrichment):
|
||
counters.enriched += 1
|
||
else:
|
||
# #3338 (та же дыра, что #3332 у domclick): карточка взята и
|
||
# разобрана, а UPDATE не задел ни одной строки — объявление
|
||
# удалено/деактивировано между снимком и записью. Попытка была,
|
||
# исхода не было: attempted переставал сходиться с
|
||
# enriched + blocked + gone + failed, и расхождение читается как
|
||
# потерянный отказ площадки. Исход failed: непрошедший UPDATE — не
|
||
# успех, не блок и не gone (снятие метит is_active=FALSE сам, по 404).
|
||
counters.failed += 1
|
||
# Печатаем ОБА идентификатора: WHERE в save ключуется по
|
||
# source_id из разобранного HTML (item_id), а не по row-id из
|
||
# снимка. При редиректе/подмене карточки строка listing_id
|
||
# существует и жива — не нашлась строка с source_id=item_id.
|
||
logger.warning(
|
||
"avito_detail_backfill: run_id=%d listing %s -- карточка "
|
||
"разобрана, но UPDATE не нашёл строку: listing_id=%s item_id=%s "
|
||
"(WHERE по item_id из HTML)",
|
||
run_id,
|
||
source_url,
|
||
row["id"],
|
||
enrichment.item_id,
|
||
)
|
||
if use_curl:
|
||
items_since_warm += 1
|
||
breaker.record_success()
|
||
consecutive_failures = 0
|
||
|
||
except AvitoListingGoneError as gone_exc:
|
||
# #2034: мёртвый листинг (404 / removed) — НЕ блок, НЕ failed.
|
||
# Координатные дыры в lat-null очереди в основном dead-листинги;
|
||
# browser-mode рендерит их «Ошибка 404» без item-view → раньше это
|
||
# ловилось как soft-block → consecutive-block breaker абортил run до
|
||
# live-листингов (run 458: attempted=5 enriched=0 blocked=5 → abort).
|
||
# 404 нейтрален к breaker'у (#3184: record_neutral -- ни окно, ни серия,
|
||
# ни safety-net не трогаются). Метим is_active=FALSE → листинг уходит из
|
||
# scope (snapshot SELECT фильтрует is_active = TRUE) и не тратит фетчи впредь.
|
||
counters.gone += 1
|
||
breaker.record_neutral()
|
||
# 404 — честный ответ площадки, значит тракт цел: серия отказов
|
||
# прерывается (блочный брейкер 404 не трогает — см. #2034).
|
||
consecutive_failures = 0
|
||
failure_census[_failure_signature(gone_exc)] += 1
|
||
try:
|
||
with db.begin_nested():
|
||
db.execute(
|
||
text("UPDATE listings SET is_active = FALSE WHERE id = :id"),
|
||
{"id": row["id"]},
|
||
)
|
||
# #2674: 404 с площадки — самый достоверный сигнал снятия,
|
||
# фиксируем его в дневной истории (listings_snapshots.status
|
||
# был константой 'active' у всех строк, 394 299). Тот же
|
||
# SAVEPOINT, что и UPDATE флага: снимок без флага (или
|
||
# наоборот) невозможен. price_rub из snapshot-SELECT —
|
||
# .get() консервативен ради mock-снапшотов старых тестов.
|
||
gone_price = row.get("price_rub")
|
||
if gone_price is not None:
|
||
upsert_listing_snapshot(
|
||
db,
|
||
listing_id=row["id"],
|
||
price_rub=gone_price,
|
||
run_id=run_id,
|
||
status="closed",
|
||
)
|
||
except Exception:
|
||
logger.warning(
|
||
"avito_detail_backfill: run_id=%d failed to mark listing %s "
|
||
"is_active=FALSE",
|
||
run_id,
|
||
row["id"],
|
||
exc_info=True,
|
||
)
|
||
logger.info(
|
||
"avito_detail_backfill: run_id=%d GONE #%d/%d listing %s -> "
|
||
"marked is_active=FALSE",
|
||
run_id,
|
||
counters.gone,
|
||
len(snapshot),
|
||
source_url,
|
||
)
|
||
|
||
except (AvitoBlockedError, AvitoRateLimitedError) as e:
|
||
# #3288: пустой пул — не блок и не отказ площадки: запрос не
|
||
# уходил вовсе, следующая карточка упрётся ровно в то же самое.
|
||
# Сюда он приезжает под видом блокировки, потому что fetch_detail
|
||
# заворачивает в AvitoSidecarUnavailableError любое исключение
|
||
# фетча. Исход честно failed (тот же разряд, что транспортные
|
||
# сбои ниже) — тождество attempted = enriched + blocked + gone +
|
||
# failed остаётся целым (#3338), а причину несёт no_proxy_stop=1
|
||
# и mark_failed с текстом про пул, как у домклика после #3283.
|
||
if _caused_by_empty_pool(e):
|
||
counters.failed += 1
|
||
logger.error(
|
||
"avito_detail_backfill: run_id=%d СТОП — пул прокси пуст, "
|
||
"к площадке не ходили. enriched=%d attempted=%d",
|
||
run_id,
|
||
counters.enriched,
|
||
counters.attempted,
|
||
)
|
||
no_proxy_stop = True
|
||
break
|
||
|
||
ban_kind = ban_kind_of_exception(e)
|
||
# #3288: в окно доли идёт только 'platform' — infra (отказ нашего
|
||
# сайдкара) уходит в знаменатель, см. BlockRatioBreaker.record_block.
|
||
breaker.record_block(ban_kind)
|
||
counters.blocked += 1
|
||
failure_census[_failure_signature(e)] += 1
|
||
block_ban_kinds[ban_kind] += 1
|
||
do_sleep = False
|
||
# #3251/#3283g: на настоящий бан ПЛОЩАДКОЙ (AvitoBlockedError и подтипы:
|
||
# AvitoContentBlockedError, AvitoWarmupCookiesMissingError; НЕ на
|
||
# AvitoSidecarUnavailableError — подтип AvitoRateLimitedError, отказ
|
||
# НАШЕГО тракта, площадка ни при чём) есть два инструмента:
|
||
# 1) rotate-on-ban (#3283g) — меняет exit-IP И сбрасывает context
|
||
# (через _rotate_current_proxy), бюджетно (rotate_on_ban_max за
|
||
# прогон) и с gap'ом (rotate_on_ban_min_gap попыток от последней
|
||
# ротации любого происхождения, gap считаем по
|
||
# attempts_since_rotation — она и так обнуляется на ротациях).
|
||
# Смена IP лечит бан вживую за секунды (см. module docstring
|
||
# задачи) и НЕ подвержена каскаду ниже — новый адрес не хранит
|
||
# сожжённый PoW старого.
|
||
# 2) bare context reset (#3251) — БЕЗ смены IP, ровно один раз за
|
||
# прогон (domclick_detail_backfill #3212 — тот же приём). Сброс
|
||
# на каждый блок сам себя поддерживает: пройденный QRATOR-PoW
|
||
# живёт в context'е, сброс его выбрасывает → следующая проверка
|
||
# с ТОГО ЖЕ IP снова блокируется — необратимый каскад.
|
||
# Приоритет — rotate-on-ban: если она сработала, context уже сброшен
|
||
# вместе со сменой IP, поэтому bare reset для этого события и всех
|
||
# последующих (пока context_reset_used не сброшен — а он не сбрасывается
|
||
# в рамках прогона) больше не нужен. Оба сброса в ОДНОМ событии не
|
||
# допускаем: разные механизмы, но request_context_reset() внутри один,
|
||
# и два подряд — то самое "два сброса подряд", которого просит избежать
|
||
# задача, даже если один из них идёт с новым IP.
|
||
if browser_fetcher is not None and isinstance(e, AvitoBlockedError):
|
||
ban_budget = settings.avito_detail_backfill_rotate_on_ban_max
|
||
if ban_budget > 0 and rotate_on_ban_used >= ban_budget:
|
||
logger.info(
|
||
"avito_detail_backfill: run_id=%d rotate-on-ban budget "
|
||
"exhausted (%d/%d) -- staying on current proxy",
|
||
run_id,
|
||
rotate_on_ban_used,
|
||
ban_budget,
|
||
)
|
||
elif (
|
||
ban_budget > 0
|
||
and attempts_since_rotation
|
||
< settings.avito_detail_backfill_rotate_on_ban_min_gap
|
||
):
|
||
logger.info(
|
||
"avito_detail_backfill: run_id=%d rotate-on-ban skipped -- "
|
||
"only %d attempts since last rotation (need >=%d)",
|
||
run_id,
|
||
attempts_since_rotation,
|
||
settings.avito_detail_backfill_rotate_on_ban_min_gap,
|
||
)
|
||
elif ban_budget > 0:
|
||
# Бюджет тратим и на неудачу — иначе отказавший rotate_proxy
|
||
# дёргался бы на КАЖДОМ следующем бане до конца прогона.
|
||
rotate_on_ban_used += 1
|
||
attempts_since_rotation = 0
|
||
logger.warning(
|
||
"avito_detail_backfill: run_id=%d rotate-on-ban #%d/%d -- "
|
||
"platform ban, rotating exit IP instead of bare reset",
|
||
run_id,
|
||
rotate_on_ban_used,
|
||
ban_budget,
|
||
)
|
||
# context_reset_used выставляем ТОЛЬКО по факту успеха: при
|
||
# тихом отказе rotate_proxy (лимит, нет rotate_url, сеть)
|
||
# сброса контекста внутри НЕ произошло, и съесть им
|
||
# одноразовый bare-reset #3251 значило бы остаться и без
|
||
# нового IP, и без сброса вообще.
|
||
if await _rotate_current_proxy(
|
||
db, run_id, browser_fetcher, trigger="rotate-on-ban"
|
||
):
|
||
context_reset_used = True
|
||
|
||
if not context_reset_used:
|
||
context_reset_used = True
|
||
browser_fetcher.request_context_reset()
|
||
logger.warning(
|
||
"avito_detail_backfill: run_id=%d BLOCKED #%d/%d (consecutive=%d): %s",
|
||
run_id,
|
||
idx + 1,
|
||
len(snapshot),
|
||
breaker.consecutive_blocks,
|
||
e,
|
||
)
|
||
# abort FIRST — на последнем блоке не тратим cooldown+research/rotate.
|
||
# #3184: доля в скользящем окне ИЛИ safety-net (все попытки — блоки, ни
|
||
# одного успеха, накопилось max_consecutive_blocks) -- см. module docstring.
|
||
abort_reason = breaker.abort_reason()
|
||
if abort_reason is not None:
|
||
# Печатаем ту величину, по которой обрыв и произошёл. Раньше здесь
|
||
# всегда стояло "%d consecutive blocks", хотя критерий с #3184 стал
|
||
# ratio: прогон 5210 (14 блоков из 20, ровно порог) отпечатал
|
||
# "ABORT -- 1 consecutive blocks" -- текущая серия в тот момент
|
||
# действительно равнялась единице, но обрыв был не по ней.
|
||
logger.error(
|
||
"avito_detail_backfill: run_id=%d ABORT -- %s, "
|
||
"частая причина: %s. enriched=%d attempted=%d",
|
||
run_id,
|
||
breaker.abort_explanation(),
|
||
_top_failure(failure_census) or "причина не определена",
|
||
counters.enriched,
|
||
counters.attempted,
|
||
)
|
||
aborted_by_blocks = True
|
||
break
|
||
# МГТС sticky-IP: один фикс. exit-IP, per-connection ротации нет (проверено:
|
||
# 6/6 свежих сессий = тот же IP 109.252.125.80; ротация только вручную
|
||
# кнопкой). Уйти на свежий IP софтом нельзя → блок = rate-limit текущего IP:
|
||
# даём окну остыть (cooldown) и освежаем куки in-session, НЕ дискардим
|
||
# рабочую прогретую сессию (rebuild на том же IP бесполезен для escape +
|
||
# грузит rate-limited IP полным прогревом). items_since_warm НЕ сбрасываем.
|
||
if use_curl:
|
||
await asyncio.sleep(block_cooldown_sec)
|
||
if session is not None:
|
||
# #1950: bound сетевой re-search.
|
||
try:
|
||
await asyncio.wait_for(
|
||
research_in_session(session), timeout=warm_timeout_s
|
||
)
|
||
except Exception:
|
||
logger.warning(
|
||
"avito_detail_backfill: run_id=%d on-block re-search "
|
||
"timed out/failed",
|
||
run_id,
|
||
exc_info=True,
|
||
)
|
||
else:
|
||
# #1950: bound rotate — _rotate_ip ждёт settle + до 3 changeip-попыток;
|
||
# на блокирующем changeip (зависшее соединение) без timeout loop виснет.
|
||
try:
|
||
await asyncio.wait_for(scraper._rotate_ip(), timeout=rotate_timeout_s)
|
||
except Exception:
|
||
logger.warning(
|
||
"avito_detail_backfill: run_id=%d rotate_ip timed out/failed "
|
||
"(>%.0fs) -- continuing",
|
||
run_id,
|
||
rotate_timeout_s,
|
||
exc_info=True,
|
||
)
|
||
|
||
except TimeoutError as e:
|
||
# asyncio.wait_for → TimeoutError (py3.12: asyncio.TimeoutError — alias).
|
||
# Ловим ДО общего Exception (TimeoutError ⊂ OSError ⊂ Exception). Зависший
|
||
# fetch отменён → листинг failed, переходим к следующему (loop не зависает,
|
||
# run не zombie #1950). Не считаем soft-блоком: rotate не дёргаем.
|
||
counters.failed += 1
|
||
consecutive_failures += 1
|
||
breaker.record_failure()
|
||
failure_census[_failure_signature(e)] += 1
|
||
logger.warning(
|
||
"avito_detail_backfill: run_id=%d listing %s TIMEOUT (>%.0fs) -- skip",
|
||
run_id,
|
||
source_url,
|
||
fetch_timeout_s,
|
||
)
|
||
try:
|
||
db.rollback()
|
||
except Exception:
|
||
pass
|
||
|
||
except Exception as e:
|
||
counters.failed += 1
|
||
consecutive_failures += 1
|
||
breaker.record_failure()
|
||
failure_census[_failure_signature(e)] += 1
|
||
logger.warning(
|
||
"avito_detail_backfill: run_id=%d listing %s failed: %s",
|
||
run_id,
|
||
source_url,
|
||
e,
|
||
)
|
||
try:
|
||
db.rollback()
|
||
except Exception:
|
||
pass
|
||
|
||
# Ротация exit-IP по счётчику попыток на арендованном прокси (задача поверх
|
||
# #3298). browser_fetcher is not None -- см. _rotate_current_proxy docstring
|
||
# за тем, почему только browser-режим держит lease с proxy_id. Обнуляем
|
||
# счётчик ДО вызова (не после) -- неудача ротации не должна дёргать
|
||
# rotate_proxy на КАЖДОЙ следующей попытке до конца прогона.
|
||
if (
|
||
browser_fetcher is not None
|
||
and attempts_since_rotation >= settings.avito_detail_backfill_rotate_after_attempts
|
||
):
|
||
attempts_since_rotation = 0
|
||
await _rotate_current_proxy(db, run_id, browser_fetcher)
|
||
|
||
if consecutive_failures >= max_consecutive_failures:
|
||
logger.error(
|
||
"avito_detail_backfill: run_id=%d ABORT -- %d отказов подряд без "
|
||
"единого успеха, частая причина: %s. enriched=%d attempted=%d",
|
||
run_id,
|
||
consecutive_failures,
|
||
_top_failure(failure_census) or "неизвестна",
|
||
counters.enriched,
|
||
counters.attempted,
|
||
)
|
||
break
|
||
|
||
if counters.attempted % 25 == 0:
|
||
current_counters = counters.to_dict()
|
||
runs_mod.update_heartbeat(db, run_id, current_counters)
|
||
|
||
counters.duration_sec = time.monotonic() - start
|
||
current_counters = counters.to_dict()
|
||
# #3184: гистограмма длин пачек блоков (streak -> сколько раз встретилась) --
|
||
# иначе эффект правки на #2674-статистике нечем измерить постфактум. finalize()
|
||
# досчитывает хвостовую пачку, если прогон оборвался посреди серии.
|
||
streak_histogram = breaker.finalize()
|
||
if streak_histogram:
|
||
current_counters["block_streak_histogram"] = streak_histogram # type: ignore[assignment]
|
||
# Какой из двух критериев оборвал прогон -- в counters, а не только в логах:
|
||
# разбор простоя идёт SQL-запросом по scrape_runs, а не грепом контейнера,
|
||
# и без этого ключа "banned" опять не отличить по причине (#3178).
|
||
if abort_reason is not None:
|
||
current_counters["abort_reason"] = abort_reason # type: ignore[assignment]
|
||
if no_proxy_stop:
|
||
# #3288 (как #3283 у домклика): остановка из-за пустого пула — НЕ блок,
|
||
# поэтому и не aborted_by_blocks: иначе прогон уйдёт в 'banned' и запись
|
||
# будет утверждать про площадку то, чего не было. Это отказ нашей стороны.
|
||
current_counters["no_proxy_stop"] = 1
|
||
runs_mod.mark_failed(
|
||
db,
|
||
run_id,
|
||
"пул прокси пуст — к площадке не ходили (#3288)",
|
||
current_counters,
|
||
)
|
||
logger.info(
|
||
"avito_detail_backfill: run_id=%d FINISHED (пул пуст) -- attempted=%d "
|
||
"enriched=%d blocked=%d gone=%d failed=%d duration=%.1fs",
|
||
run_id,
|
||
counters.attempted,
|
||
counters.enriched,
|
||
counters.blocked,
|
||
counters.gone,
|
||
counters.failed,
|
||
counters.duration_sec,
|
||
)
|
||
return counters
|
||
runs_mod.mark_backfill_finished(
|
||
db,
|
||
run_id,
|
||
current_counters,
|
||
source="avito_detail_backfill",
|
||
aborted_by_blocks=aborted_by_blocks,
|
||
fail_hint=_top_failure(failure_census),
|
||
ban_kinds=block_ban_kinds,
|
||
)
|
||
logger.info(
|
||
"avito_detail_backfill: run_id=%d FINISHED -- attempted=%d enriched=%d "
|
||
"blocked=%d gone=%d failed=%d duration=%.1fs",
|
||
run_id,
|
||
counters.attempted,
|
||
counters.enriched,
|
||
counters.blocked,
|
||
counters.gone,
|
||
counters.failed,
|
||
counters.duration_sec,
|
||
)
|
||
return counters
|
||
|
||
except Exception as exc:
|
||
counters.duration_sec = time.monotonic() - start
|
||
logger.exception(
|
||
"avito_detail_backfill: run_id=%d FAILED after %.1fs",
|
||
run_id,
|
||
counters.duration_sec,
|
||
)
|
||
runs_mod.mark_failed(db, run_id, str(exc)[:1000], counters.to_dict())
|
||
raise
|
||
|
||
finally:
|
||
if own_session and session is not None:
|
||
await session.close()
|
||
if own_browser and browser_fetcher is not None:
|
||
await browser_fetcher.__aexit__(None, None, None)
|