Compare commits
No commits in common. "17548e90460cc42a34a296ff2dbd5585721d4194" and "a941d1899f83fbd7c2a709f6bdc8bc2afb81a531" have entirely different histories.
17548e9046
...
a941d1899f
3 changed files with 1 additions and 235 deletions
|
|
@ -1,175 +0,0 @@
|
||||||
"""Чекпоинт по якорям для avito_city_sweep (#3074).
|
|
||||||
|
|
||||||
Прод-факт, из которого выросла задача. Типичный итог свипа:
|
|
||||||
|
|
||||||
{"anchors_done": 1, "anchors_total": 5, ..., "enrichment_abort_note":
|
|
||||||
"detail enrichment aborted (Avito detail firewall/soft-block ...)"}
|
|
||||||
|
|
||||||
То есть прогон срывается блокировкой на ПЕРВОМ из пяти якорей — 25 банов за
|
|
||||||
60 дней. Без чекпоинта следующий прогон снова начинает с первого якоря,
|
|
||||||
упирается в ту же стену, и якоря 2-5 не собираются никогда.
|
|
||||||
|
|
||||||
Ключ чекпоинта — ИМЯ якоря, а не его индекс: состав списка зависит от
|
|
||||||
`city_slug`, и позиция в нём не устойчива между городами.
|
|
||||||
|
|
||||||
ИНВАРИАНТ, РАДИ КОТОРОГО ТЕСТ. В чекпоинт попадает только якорь, пройденный до
|
|
||||||
конца. Ветка блокировки делает `return` и до записи не доходит, а вот
|
|
||||||
`except Exception` — доходит: якорь упал, но цикл продолжается. Записать такой
|
|
||||||
якорь как пройденный значило бы, что следующий прогон его пропустит и
|
|
||||||
объявления оттуда не соберутся НИКОГДА, причём молча — прогон завершится
|
|
||||||
штатно. Ровно та же граница, что у combo в yandex-свипе.
|
|
||||||
"""
|
|
||||||
|
|
||||||
from __future__ import annotations
|
|
||||||
|
|
||||||
import os
|
|
||||||
|
|
||||||
# Settings собирается автофикстурой conftest'а и требует database_url. Выставляем
|
|
||||||
# до остальных импортов — так же, как в test_3074_yandex_sweep_checkpoint.py.
|
|
||||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
|
||||||
|
|
||||||
import json
|
|
||||||
import types
|
|
||||||
from typing import Any
|
|
||||||
from unittest.mock import MagicMock, patch
|
|
||||||
|
|
||||||
import pytest
|
|
||||||
|
|
||||||
ANCHOR_A = (56.83, 60.60, "ekb-center")
|
|
||||||
ANCHOR_B = (56.79, 60.63, "ekb-south")
|
|
||||||
|
|
||||||
|
|
||||||
class _FakeDb:
|
|
||||||
"""Резюм-SELECT отдаёт counters предшественника; heartbeat'ы записываются."""
|
|
||||||
|
|
||||||
def __init__(self, prev_counters: dict[str, Any] | None = None) -> None:
|
|
||||||
self.prev_counters = prev_counters or {}
|
|
||||||
self.heartbeats: list[dict[str, Any]] = []
|
|
||||||
|
|
||||||
def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> Any:
|
|
||||||
if params and "counters" in params:
|
|
||||||
self.heartbeats.append(json.loads(params["counters"]))
|
|
||||||
return MagicMock()
|
|
||||||
if params and "rid" in params:
|
|
||||||
return MagicMock(
|
|
||||||
fetchone=lambda: types.SimpleNamespace(counters=self.prev_counters)
|
|
||||||
)
|
|
||||||
return MagicMock()
|
|
||||||
|
|
||||||
def commit(self) -> None: ...
|
|
||||||
def rollback(self) -> None: ...
|
|
||||||
|
|
||||||
|
|
||||||
class _FakeAsyncSession:
|
|
||||||
def __init__(self, *_a: Any, **_kw: Any) -> None: ...
|
|
||||||
async def __aenter__(self) -> _FakeAsyncSession:
|
|
||||||
return self
|
|
||||||
|
|
||||||
async def __aexit__(self, *_e: Any) -> None:
|
|
||||||
return None
|
|
||||||
|
|
||||||
|
|
||||||
class _FakeScraper:
|
|
||||||
"""Двойник AvitoScraper: помнит, за какими якорями реально ходили."""
|
|
||||||
|
|
||||||
visited: list[tuple[float, float]] = [] # noqa: RUF012 — тестовый сборник
|
|
||||||
raise_on: tuple[float, float] | None = None
|
|
||||||
|
|
||||||
def __init__(self, *_a: Any, **_kw: Any) -> None:
|
|
||||||
self._browser = None
|
|
||||||
self._cffi = None
|
|
||||||
|
|
||||||
async def fetch_around(self, lat: float, lon: float, *_a: Any, **_kw: Any) -> list:
|
|
||||||
_FakeScraper.visited.append((lat, lon))
|
|
||||||
if _FakeScraper.raise_on == (lat, lon):
|
|
||||||
raise RuntimeError("якорь упал по не-баново́й причине")
|
|
||||||
return []
|
|
||||||
|
|
||||||
|
|
||||||
def _config() -> types.SimpleNamespace:
|
|
||||||
return types.SimpleNamespace(
|
|
||||||
scraper_fetch_mode="cffi",
|
|
||||||
scraper_proxy_url=None,
|
|
||||||
use_proxy_pool_browser=False,
|
|
||||||
browser_http_endpoint=None,
|
|
||||||
environment="test",
|
|
||||||
avito_serp_ok_not_banned=True,
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
async def _run(prev: dict[str, Any] | None, raise_on: tuple[float, float] | None = None):
|
|
||||||
from scraper_kit.orchestration import pipeline as pl
|
|
||||||
|
|
||||||
_FakeScraper.visited = []
|
|
||||||
_FakeScraper.raise_on = raise_on
|
|
||||||
db = _FakeDb(prev)
|
|
||||||
|
|
||||||
with (
|
|
||||||
patch.object(pl, "AvitoScraper", _FakeScraper),
|
|
||||||
patch.object(pl, "AsyncSession", _FakeAsyncSession),
|
|
||||||
patch.object(pl, "save_listings", lambda *_a, **_kw: (0, 0)),
|
|
||||||
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
|
|
||||||
):
|
|
||||||
await pl.run_avito_city_sweep(
|
|
||||||
db, # type: ignore[arg-type]
|
|
||||||
run_id=7001,
|
|
||||||
config=_config(),
|
|
||||||
matcher=MagicMock(),
|
|
||||||
enrichment=MagicMock(),
|
|
||||||
anchors=[ANCHOR_A, ANCHOR_B],
|
|
||||||
enrich_houses=False,
|
|
||||||
enrich_imv=False,
|
|
||||||
detail_top_n=0,
|
|
||||||
resume_run_id=6999 if prev is not None else None,
|
|
||||||
)
|
|
||||||
return db
|
|
||||||
|
|
||||||
|
|
||||||
def _last_checkpoint(db: _FakeDb) -> list[str]:
|
|
||||||
with_ckpt = [hb for hb in db.heartbeats if "done_buckets" in hb]
|
|
||||||
assert with_ckpt, "ни один heartbeat не унёс done_buckets — чекпоинт не персистится"
|
|
||||||
return with_ckpt[-1]["done_buckets"]
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
|
||||||
async def test_checkpointed_anchor_is_skipped_without_a_single_request() -> None:
|
|
||||||
"""Якорь из чекпоинта не опрашивается вовсе — ни одного обращения к источнику.
|
|
||||||
|
|
||||||
Ядро задачи: до фикса повторный прогон снова шёл в первый якорь и снова
|
|
||||||
получал там бан.
|
|
||||||
"""
|
|
||||||
await _run({"done_buckets": ["ekb-center"]})
|
|
||||||
|
|
||||||
assert (ANCHOR_A[0], ANCHOR_A[1]) not in _FakeScraper.visited, (
|
|
||||||
"якорь из чекпоинта всё-таки опрашивали"
|
|
||||||
)
|
|
||||||
assert (ANCHOR_B[0], ANCHOR_B[1]) in _FakeScraper.visited, "второй якорь не обошли"
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
|
||||||
async def test_checkpoint_accumulates_over_inherited() -> None:
|
|
||||||
"""Пройденный якорь дописывается поверх унаследованных, а не затирает их."""
|
|
||||||
db = await _run({"done_buckets": ["ekb-center"]})
|
|
||||||
assert _last_checkpoint(db) == ["ekb-center", "ekb-south"]
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
|
||||||
async def test_without_resume_all_anchors_are_visited() -> None:
|
|
||||||
"""Без чекпоинта поведение прежнее — обходятся все якоря."""
|
|
||||||
db = await _run(None)
|
|
||||||
assert len(_FakeScraper.visited) == 2
|
|
||||||
assert _last_checkpoint(db) == ["ekb-center", "ekb-south"]
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
|
||||||
async def test_failed_anchor_does_not_enter_checkpoint() -> None:
|
|
||||||
"""Упавший якорь НЕ считается пройденным.
|
|
||||||
|
|
||||||
Иначе следующий прогон пропустит его навсегда, и это будет незаметно:
|
|
||||||
прогон завершается штатно, просто часть города не собирается никогда.
|
|
||||||
"""
|
|
||||||
db = await _run(None, raise_on=(ANCHOR_A[0], ANCHOR_A[1]))
|
|
||||||
|
|
||||||
ckpt = _last_checkpoint(db)
|
|
||||||
assert "ekb-center" not in ckpt, "упавший якорь попал в чекпоинт"
|
|
||||||
assert "ekb-south" in ckpt, "исправный якорь не зафиксирован"
|
|
||||||
|
|
@ -1132,7 +1132,6 @@ async def run_avito_city_sweep(
|
||||||
request_delay_sec: float = 7.0,
|
request_delay_sec: float = 7.0,
|
||||||
enrich_imv: bool = True,
|
enrich_imv: bool = True,
|
||||||
region_code: int = DEFAULT_REGION_CODE,
|
region_code: int = DEFAULT_REGION_CODE,
|
||||||
resume_run_id: int | None = None,
|
|
||||||
) -> CitySweepCounters:
|
) -> CitySweepCounters:
|
||||||
"""Full city sweep: iterate anchors × pages → save → enrich houses + detail → IMV.
|
"""Full city sweep: iterate anchors × pages → save → enrich houses + detail → IMV.
|
||||||
|
|
||||||
|
|
@ -1155,32 +1154,6 @@ async def run_avito_city_sweep(
|
||||||
# а avito_slug (#12) — в путь URL, скоупя сам запрос на город-цель вместо ЕКБ.
|
# а avito_slug (#12) — в путь URL, скоупя сам запрос на город-цель вместо ЕКБ.
|
||||||
# None → ЕКБ-дефолт (совпадает с anchors=EKB_ANCHORS fallback ниже).
|
# None → ЕКБ-дефолт (совпадает с anchors=EKB_ANCHORS fallback ниже).
|
||||||
_anchors = anchors if anchors is not None else EKB_ANCHORS
|
_anchors = anchors if anchors is not None else EKB_ANCHORS
|
||||||
|
|
||||||
# ── Checkpoint/resume (#3074): единица обхода — ЯКОРЬ ──────────────────────
|
|
||||||
# Прод-факт: у avito_city_sweep типичный итог `anchors_done: 1` из
|
|
||||||
# `anchors_total: 5` — прогон срывается блокировкой детализации на первом же
|
|
||||||
# якоре (25 банов за 60 дней). Без чекпоинта следующий прогон снова начинает
|
|
||||||
# с первого якоря, упирается в ту же стену, и якоря 2-5 не собираются никогда.
|
|
||||||
#
|
|
||||||
# Ключ чекпоинта — ИМЯ якоря, а не индекс: список якорей зависит от
|
|
||||||
# `city_slug`, и позиция в нём не устойчива между прогонами разных городов.
|
|
||||||
_skip_anchors: set[str] = set()
|
|
||||||
if resume_run_id is not None:
|
|
||||||
_prev = db.execute(
|
|
||||||
text("SELECT counters FROM scrape_runs WHERE id = CAST(:rid AS bigint)"),
|
|
||||||
{"rid": resume_run_id},
|
|
||||||
).fetchone()
|
|
||||||
if _prev is not None and _prev.counters:
|
|
||||||
_pc: dict = _prev.counters if isinstance(_prev.counters, dict) else {}
|
|
||||||
_skip_anchors = set(_pc.get("done_buckets", []))
|
|
||||||
logger.info(
|
|
||||||
"city-sweep run_id=%d: resuming from run %s — %d якорей уже пройдено",
|
|
||||||
run_id,
|
|
||||||
resume_run_id,
|
|
||||||
len(_skip_anchors),
|
|
||||||
)
|
|
||||||
_done_anchors: set[str] = set(_skip_anchors)
|
|
||||||
|
|
||||||
_loc = get_city_location(city_slug)
|
_loc = get_city_location(city_slug)
|
||||||
# #262 wave 2: avito_slug у CityLocation Optional — не у каждого известного города
|
# #262 wave 2: avito_slug у CityLocation Optional — не у каждого известного города
|
||||||
# он подтверждён (403/429 на исчерпанном пуле при проверке, либо omonym-коллизия).
|
# он подтверждён (403/429 на исчерпанном пуле при проверке, либо omonym-коллизия).
|
||||||
|
|
@ -1264,19 +1237,6 @@ async def run_avito_city_sweep(
|
||||||
)
|
)
|
||||||
try:
|
try:
|
||||||
for idx, (lat, lon, name) in enumerate(_anchors, start=1):
|
for idx, (lat, lon, name) in enumerate(_anchors, start=1):
|
||||||
if name in _skip_anchors:
|
|
||||||
# #3074: якорь собран предыдущим оборванным прогоном — ни одного
|
|
||||||
# HTTP-запроса. Счётчик двигаем, чтобы `anchors_done` продолжал
|
|
||||||
# означать «докуда дошли по списку», а не «сколько собрал этот run».
|
|
||||||
counters.anchors_done = idx
|
|
||||||
logger.info(
|
|
||||||
"city-sweep run_id=%d: anchor #%d/%d (%s) пропущен — есть в чекпоинте",
|
|
||||||
run_id,
|
|
||||||
idx,
|
|
||||||
len(_anchors),
|
|
||||||
name,
|
|
||||||
)
|
|
||||||
continue
|
|
||||||
if runs.is_cancelled(db, run_id):
|
if runs.is_cancelled(db, run_id):
|
||||||
logger.info(
|
logger.info(
|
||||||
"city-sweep run_id=%d: cancelled at anchor #%d/%d (%s)",
|
"city-sweep run_id=%d: cancelled at anchor #%d/%d (%s)",
|
||||||
|
|
@ -1313,10 +1273,6 @@ async def run_avito_city_sweep(
|
||||||
lon,
|
lon,
|
||||||
)
|
)
|
||||||
|
|
||||||
# #3074: сбрасывается на КАЖДОЙ итерации — иначе один упавший якорь
|
|
||||||
# заразил бы все последующие, и чекпоинт не пополнялся бы вовсе.
|
|
||||||
_anchor_ok = True
|
|
||||||
|
|
||||||
# Capture loop variables in default args (B023): prevents stale binding
|
# Capture loop variables in default args (B023): prevents stale binding
|
||||||
# if the coroutine is scheduled after the loop variable changes.
|
# if the coroutine is scheduled after the loop variable changes.
|
||||||
_a_lat, _a_lon, _a_name = lat, lon, name
|
_a_lat, _a_lon, _a_name = lat, lon, name
|
||||||
|
|
@ -1811,20 +1767,9 @@ async def run_avito_city_sweep(
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.exception("city-sweep run_id=%d: anchor %s failed", run_id, name)
|
logger.exception("city-sweep run_id=%d: anchor %s failed", run_id, name)
|
||||||
counters.errors_count += 1
|
counters.errors_count += 1
|
||||||
_anchor_ok = False
|
|
||||||
|
|
||||||
counters.anchors_done = idx
|
counters.anchors_done = idx
|
||||||
# #3074: в чекпоинт попадает ТОЛЬКО якорь, пройденный до конца.
|
runs.update_heartbeat(db, run_id, counters.to_dict())
|
||||||
# Ветка блокировки выше делает `return` и сюда не доходит, а вот
|
|
||||||
# generic-except доходит — якорь упал, но цикл продолжается. Записать
|
|
||||||
# его как пройденный значило бы, что следующий прогон его пропустит и
|
|
||||||
# объявления оттуда не соберутся НИКОГДА, причём молча: прогон
|
|
||||||
# завершится штатно. Тот же инвариант, что у combo в yandex-свипе.
|
|
||||||
if _anchor_ok:
|
|
||||||
_done_anchors.add(name)
|
|
||||||
runs.update_heartbeat(
|
|
||||||
db, run_id, {**counters.to_dict(), "done_buckets": sorted(_done_anchors)}
|
|
||||||
)
|
|
||||||
|
|
||||||
# ── IMV-фаза: финальный обход тронутых домов ──────────
|
# ── IMV-фаза: финальный обход тронутых домов ──────────
|
||||||
if enrich_imv and all_touched_house_ids:
|
if enrich_imv and all_touched_house_ids:
|
||||||
|
|
|
||||||
|
|
@ -786,10 +786,6 @@ async def _job_avito_city_sweep(
|
||||||
request_delay_sec=float(params.get("request_delay_sec", 7.0)),
|
request_delay_sec=float(params.get("request_delay_sec", 7.0)),
|
||||||
enrich_houses=bool(params.get("enrich_houses", True)),
|
enrich_houses=bool(params.get("enrich_houses", True)),
|
||||||
radius_m=int(params.get("radius_m", 1500)),
|
radius_m=int(params.get("radius_m", 1500)),
|
||||||
# #3074: подхват якорей у оборванного предшественника. Ключ чекпоинта —
|
|
||||||
# имя якоря, оно не зависит от количества якорей, поэтому в отличие от
|
|
||||||
# combo-чекпоинта yandex-свипа гарда по числу якорей здесь не требуется.
|
|
||||||
resume_run_id=_pick_resume(db, run_id),
|
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue