Свип ДомКлика сохраняет лоты по корзинам, а не одним махом в конце (#3592)
All checks were successful
Deploy Trade-In / changes (push) Successful in 17s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 4m4s
Deploy Trade-In / build-backend (push) Successful in 2m23s
Deploy Trade-In / deploy (push) Successful in 8m6s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / perimeter-smoke (push) Successful in 1m43s
All checks were successful
Deploy Trade-In / changes (push) Successful in 17s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 4m4s
Deploy Trade-In / build-backend (push) Successful in 2m23s
Deploy Trade-In / deploy (push) Successful in 8m6s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / perimeter-smoke (push) Successful in 1m43s
This commit is contained in:
parent
98f587add6
commit
e5212e52ab
8 changed files with 419 additions and 47 deletions
|
|
@ -544,6 +544,9 @@ async def _job_domclick_city_sweep(
|
||||||
region_code=kit_resolve_region_code(params),
|
region_code=kit_resolve_region_code(params),
|
||||||
resume_run_id=kit_pick_resume(db, run_id),
|
resume_run_id=kit_pick_resume(db, run_id),
|
||||||
cookies=cookies,
|
cookies=cookies,
|
||||||
|
watchdog_sec=(
|
||||||
|
int(params["watchdog_sec"]) if params.get("watchdog_sec") is not None else None
|
||||||
|
),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -173,10 +173,14 @@ async def test_anchor_timeout_skips_anchor_and_continues(kind: str) -> None:
|
||||||
async def test_domclick_unsaved_buckets_are_not_checkpointed(broken: str) -> None:
|
async def test_domclick_unsaved_buckets_are_not_checkpointed(broken: str) -> None:
|
||||||
"""Корзина без сохранённых строк не попадает в чекпоинт.
|
"""Корзина без сохранённых строк не попадает в чекпоинт.
|
||||||
|
|
||||||
Лоты Домклика копятся в памяти и пишутся ОДНИМ save_listings после всех корзин.
|
С инкрементальным сохранением (fix/domclick-incremental-save) save_listings зовётся из
|
||||||
Фаза, снятая watchdog'ом (или упавшая на save), не сохранила ничего — даже из
|
колбэка on_bucket ПОСЛЕ каждой корзины, а чекпоинт пишется там же и только
|
||||||
корзины 'st', которую скрейпер успел пройти. Отметить её пройденной значило бы,
|
ПОСЛЕ успешного save. Поэтому корзина, чей save упал, в done_buckets не попадает,
|
||||||
что следующий прогон пропустит её через skip_buckets навсегда (миграция 308).
|
хотя скрейпер её фетч прошёл (_s.completed_buckets её содержит). Снятая
|
||||||
|
watchdog'ом фаза — тот же инвариант: до on_bucket она не дошла.
|
||||||
|
|
||||||
|
Отметить такую корзину пройденной значило бы, что следующий прогон пропустит её
|
||||||
|
через skip_buckets навсегда (миграция 308).
|
||||||
"""
|
"""
|
||||||
|
|
||||||
class _Dc(_Scraper):
|
class _Dc(_Scraper):
|
||||||
|
|
@ -185,10 +189,17 @@ async def test_domclick_unsaved_buckets_are_not_checkpointed(broken: str) -> Non
|
||||||
buckets_completed, buckets_total = 1, 6
|
buckets_completed, buckets_total = 1, 6
|
||||||
completed_buckets = ["st"] # noqa: RUF012
|
completed_buckets = ["st"] # noqa: RUF012
|
||||||
|
|
||||||
async def fetch_city(self, **_kw: Any) -> list[Any]:
|
async def fetch_city(self, **kw: Any) -> list[Any]:
|
||||||
if broken == "fetch_timeout":
|
if broken == "fetch_timeout":
|
||||||
raise TimeoutError
|
raise TimeoutError
|
||||||
return [object()]
|
# Заглушка обязана ВЫЗВАТЬ колбэк: путь сохранения переехал внутрь цикла
|
||||||
|
# по корзинам (serp.py fetch_city), и стаб, который просто возвращает лоты,
|
||||||
|
# проверял бы мёртвую ветку — save_listings не был бы вызван вовсе.
|
||||||
|
on_bucket = kw.get("on_bucket")
|
||||||
|
lots = [object()]
|
||||||
|
if on_bucket is not None:
|
||||||
|
on_bucket("st", lots)
|
||||||
|
return lots
|
||||||
|
|
||||||
saved: list[int] = []
|
saved: list[int] = []
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -153,7 +153,15 @@ class _CleanSweepScraper:
|
||||||
async def __aexit__(self, *_e: Any) -> None:
|
async def __aexit__(self, *_e: Any) -> None:
|
||||||
return None
|
return None
|
||||||
|
|
||||||
async def fetch_city(self, **_kw: Any) -> list[Any]:
|
async def fetch_city(self, **kw: Any) -> list[Any]:
|
||||||
|
# Чекпоинт с fix/domclick-incremental-save набирается ИЗ колбэка on_bucket, а не
|
||||||
|
# мержится в конце из completed_buckets — стаб обязан его вызвать, иначе
|
||||||
|
# проверялась бы мёртвая ветка. Лотов нет: корзина пройдена, но пустая —
|
||||||
|
# это законный случай, save_listings для неё не зовётся, чекпоинт пишется.
|
||||||
|
on_bucket = kw.get("on_bucket")
|
||||||
|
if on_bucket is not None:
|
||||||
|
for bucket in self.completed_buckets:
|
||||||
|
on_bucket(bucket, [])
|
||||||
return []
|
return []
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
229
tradein-mvp/backend/tests/test_domclick_incremental_save.py
Normal file
229
tradein-mvp/backend/tests/test_domclick_incremental_save.py
Normal file
|
|
@ -0,0 +1,229 @@
|
||||||
|
"""fix/domclick-incremental-save: большая выдача (Москва ≈23 690 лотов против ЕКБ
|
||||||
|
≈6 300) теряла ВСЁ собранное при снятии свипа по watchdog — сохранение было ОДНО,
|
||||||
|
в самом конце `run_domclick_city_sweep`. Прод run 7344 (`domclick_city_sweep_moskva`):
|
||||||
|
за 2ч watchdog'а (11 100с) пройдено 2 бакета из 6.
|
||||||
|
|
||||||
|
Три дефекта, три слоя тестов ниже:
|
||||||
|
1. Инкрементальное сохранение по бакетам (`fetch_city.on_bucket`) — секция 1.
|
||||||
|
2. Чекпоинт врал ("пройдено" без "сохранено") — секция 1, тест
|
||||||
|
`test_on_bucket_exception_interrupts_bucket_loop` документирует нюанс, на
|
||||||
|
котором строится фикс в pipeline.py: scraper.completed_buckets отмечает бакет
|
||||||
|
как "фетч прошёл" ДО вызова on_bucket, поэтому пайплайн больше не берёт
|
||||||
|
чекпоинт оттуда — только из факта успешного on_bucket (см. комментарий у
|
||||||
|
"Перенести счётчики" в run_domclick_city_sweep).
|
||||||
|
3. Watchdog не знал о размере выдачи — секция 2, `watchdog_sec` override.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import os
|
||||||
|
from types import SimpleNamespace
|
||||||
|
from typing import Any
|
||||||
|
from unittest.mock import patch
|
||||||
|
|
||||||
|
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
from scraper_kit.orchestration.pipeline import run_domclick_city_sweep
|
||||||
|
from scraper_kit.providers.domclick.serp import (
|
||||||
|
ROOM_BUCKETS,
|
||||||
|
DomClickBlockedError,
|
||||||
|
DomClickScraper,
|
||||||
|
)
|
||||||
|
|
||||||
|
PFX = "scraper_kit.orchestration.pipeline"
|
||||||
|
|
||||||
|
|
||||||
|
# ── Секция 1: DomClickScraper.fetch_city(on_bucket=...) ──────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
def _scraper() -> DomClickScraper:
|
||||||
|
return DomClickScraper(SimpleNamespace(scraper_proxy_url=None))
|
||||||
|
|
||||||
|
|
||||||
|
class _FakeFetcherCtx:
|
||||||
|
async def __aenter__(self) -> SimpleNamespace:
|
||||||
|
return SimpleNamespace(report_ban=lambda *_a, **_k: None)
|
||||||
|
|
||||||
|
async def __aexit__(self, *_exc: Any) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
async def _run_fetch_city(
|
||||||
|
scraper: DomClickScraper,
|
||||||
|
*,
|
||||||
|
sweep_bucket: Any,
|
||||||
|
on_bucket: Any = None,
|
||||||
|
start: int = 0,
|
||||||
|
skip: set[str] | None = None,
|
||||||
|
) -> list[str]:
|
||||||
|
with (
|
||||||
|
patch.object(scraper, "_sweep_bucket", sweep_bucket),
|
||||||
|
patch(
|
||||||
|
"scraper_kit.providers._base.build_browser_fetcher",
|
||||||
|
lambda *_a, **_k: _FakeFetcherCtx(),
|
||||||
|
),
|
||||||
|
):
|
||||||
|
return await scraper.fetch_city(
|
||||||
|
city_id=1,
|
||||||
|
pages=1,
|
||||||
|
start_bucket_index=start,
|
||||||
|
skip_buckets=skip,
|
||||||
|
on_bucket=on_bucket,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
async def test_on_bucket_called_once_per_bucket_with_only_its_own_lots() -> None:
|
||||||
|
"""on_bucket зовётся по разу на каждый успешный бакет с лотами ИМЕННО его,
|
||||||
|
а не накопленным out_lots (главная регрессия #1 issue — было "одно сохранение
|
||||||
|
в конце")."""
|
||||||
|
|
||||||
|
async def _fake_sweep(*, rooms: str, out_lots: list[str], **_kw: Any) -> None:
|
||||||
|
out_lots.append(f"{rooms}-a")
|
||||||
|
out_lots.append(f"{rooms}-b")
|
||||||
|
|
||||||
|
calls: list[tuple[str, list[str]]] = []
|
||||||
|
s = _scraper()
|
||||||
|
result = await _run_fetch_city(
|
||||||
|
s, sweep_bucket=_fake_sweep, on_bucket=lambda n, lots: calls.append((n, list(lots)))
|
||||||
|
)
|
||||||
|
|
||||||
|
assert [c[0] for c in calls] == list(ROOM_BUCKETS)
|
||||||
|
for bucket, lots in calls:
|
||||||
|
assert lots == [f"{bucket}-a", f"{bucket}-b"], (bucket, lots, "получил чужие лоты")
|
||||||
|
assert len(result) == len(ROOM_BUCKETS) * 2
|
||||||
|
|
||||||
|
|
||||||
|
async def test_on_bucket_skips_failed_and_blocked_buckets() -> None:
|
||||||
|
"""Битый бакет (generic Exception) и бакет с QRATOR-блоком не отдают лоты в
|
||||||
|
on_bucket — там ничего не собрано/не гарантированно собрано."""
|
||||||
|
|
||||||
|
async def _fake_sweep_fail_middle(*, rooms: str, out_lots: list[str], **_kw: Any) -> None:
|
||||||
|
out_lots.append(f"{rooms}-x")
|
||||||
|
if rooms == "1":
|
||||||
|
raise ValueError("boom")
|
||||||
|
|
||||||
|
calls_a: list[str] = []
|
||||||
|
s_a = _scraper()
|
||||||
|
await _run_fetch_city(
|
||||||
|
s_a, sweep_bucket=_fake_sweep_fail_middle, on_bucket=lambda n, _lots: calls_a.append(n)
|
||||||
|
)
|
||||||
|
assert "1" not in calls_a
|
||||||
|
# continue идёт дальше — остальные бакеты всё равно получают on_bucket.
|
||||||
|
assert calls_a == [b for b in ROOM_BUCKETS if b != "1"], calls_a
|
||||||
|
|
||||||
|
async def _fake_sweep_block(*, rooms: str, out_lots: list[str], **_kw: Any) -> None:
|
||||||
|
out_lots.append(f"{rooms}-x")
|
||||||
|
if rooms == "1":
|
||||||
|
raise DomClickBlockedError("QRATOR")
|
||||||
|
|
||||||
|
calls_b: list[str] = []
|
||||||
|
s_b = _scraper()
|
||||||
|
await _run_fetch_city(
|
||||||
|
s_b, sweep_bucket=_fake_sweep_block, on_bucket=lambda n, _lots: calls_b.append(n)
|
||||||
|
)
|
||||||
|
# break останавливает обход целиком — после блока ни один бакет не пробуется.
|
||||||
|
assert calls_b == ["st"], calls_b
|
||||||
|
|
||||||
|
|
||||||
|
async def test_on_bucket_exception_interrupts_bucket_loop() -> None:
|
||||||
|
"""Исключение из on_bucket (канал кооперативной отмены) прерывает цикл по
|
||||||
|
ROOM_BUCKETS целиком — остальные бакеты не идут."""
|
||||||
|
|
||||||
|
async def _fake_sweep(*, rooms: str, out_lots: list[str], **_kw: Any) -> None:
|
||||||
|
out_lots.append(f"{rooms}-x")
|
||||||
|
|
||||||
|
visited: list[str] = []
|
||||||
|
|
||||||
|
def _on_bucket(name: str, lots: list[str]) -> None:
|
||||||
|
visited.append(name)
|
||||||
|
if name == "1":
|
||||||
|
raise RuntimeError("cancelled")
|
||||||
|
|
||||||
|
s = _scraper()
|
||||||
|
with pytest.raises(RuntimeError, match="cancelled"):
|
||||||
|
await _run_fetch_city(s, sweep_bucket=_fake_sweep, on_bucket=_on_bucket)
|
||||||
|
|
||||||
|
assert visited == ["st", "1"], visited
|
||||||
|
# Скрейпер-уровень успел зафетчить оба бакета ДО того, как on_bucket поднял
|
||||||
|
# исключение — на этом нюансе строится фикс чекпоинта в pipeline.py: pipeline
|
||||||
|
# больше не берёт done_buckets из scraper.completed_buckets, только из факта
|
||||||
|
# успешного on_bucket (см. run_domclick_city_sweep._on_bucket/_checkpoint).
|
||||||
|
assert s.completed_buckets == ["st", "1"], s.completed_buckets
|
||||||
|
|
||||||
|
|
||||||
|
async def test_on_bucket_none_preserves_return_value() -> None:
|
||||||
|
"""on_bucket=None → поведение прежнее, байт-в-байт: лоты возвращаются из
|
||||||
|
fetch_city как раньше, никаких промежуточных вызовов."""
|
||||||
|
|
||||||
|
async def _fake_sweep(*, rooms: str, out_lots: list[str], **_kw: Any) -> None:
|
||||||
|
out_lots.append(f"{rooms}-only")
|
||||||
|
|
||||||
|
s = _scraper()
|
||||||
|
result = await _run_fetch_city(s, sweep_bucket=_fake_sweep, on_bucket=None)
|
||||||
|
|
||||||
|
assert result == [f"{b}-only" for b in ROOM_BUCKETS]
|
||||||
|
assert s.completed_buckets == list(ROOM_BUCKETS)
|
||||||
|
|
||||||
|
|
||||||
|
# ── Секция 2: run_domclick_city_sweep(watchdog_sec=...) ──────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
class _NoOpRuns:
|
||||||
|
"""Достаточно методов, чтобы SERP-фаза дошла до asyncio.wait_for и честно
|
||||||
|
финализировалась после симулированного TimeoutError — значения не важны,
|
||||||
|
важен ТОЛЬКО timeout, с которым позвали wait_for."""
|
||||||
|
|
||||||
|
def is_cancelled(self, db: Any, run_id: int) -> bool:
|
||||||
|
return False
|
||||||
|
|
||||||
|
def update_heartbeat(self, db: Any, run_id: int, counters: dict[str, Any]) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
def mark_done(self, db: Any, run_id: int, counters: dict[str, Any]) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
def mark_failed(self, db: Any, run_id: int, error: str, counters: dict[str, Any]) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
def mark_banned(
|
||||||
|
self, db: Any, run_id: int, error: str, counters: dict[str, Any], **kw: Any
|
||||||
|
) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
async def _drive_and_capture_timeout(**kwargs: Any) -> float:
|
||||||
|
captured: dict[str, float] = {}
|
||||||
|
|
||||||
|
async def _fake_wait_for(coro: Any, timeout: float) -> None:
|
||||||
|
captured["timeout"] = timeout
|
||||||
|
coro.close()
|
||||||
|
raise TimeoutError()
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch(f"{PFX}.asyncio.wait_for", _fake_wait_for),
|
||||||
|
patch(f"{PFX}.runs", _NoOpRuns()),
|
||||||
|
):
|
||||||
|
await run_domclick_city_sweep(
|
||||||
|
object(), # type: ignore[arg-type]
|
||||||
|
config=SimpleNamespace(browser_http_endpoint="http://x:9000"),
|
||||||
|
matcher=object(),
|
||||||
|
run_id=1,
|
||||||
|
city_id=4,
|
||||||
|
**kwargs,
|
||||||
|
)
|
||||||
|
return captured["timeout"]
|
||||||
|
|
||||||
|
|
||||||
|
async def test_watchdog_sec_default_formula_matches_prod_11100() -> None:
|
||||||
|
"""watchdog_sec=None (дефолт) → прежняя формула байт-в-байт. pages=100,
|
||||||
|
delay=6.0 — те же параметры, что дали 11 100с в run 7344 (Москва)."""
|
||||||
|
timeout = await _drive_and_capture_timeout(pages=100, request_delay_sec=6.0)
|
||||||
|
assert timeout == 11100
|
||||||
|
|
||||||
|
|
||||||
|
async def test_watchdog_sec_override_bypasses_formula() -> None:
|
||||||
|
"""watchdog_sec задан явно → формула не считается вовсе, идёт ровно override
|
||||||
|
— даже с теми же pages/delay, что в тесте формулы выше."""
|
||||||
|
timeout = await _drive_and_capture_timeout(pages=100, request_delay_sec=6.0, watchdog_sec=777)
|
||||||
|
assert timeout == 777
|
||||||
|
|
@ -411,8 +411,19 @@ async def _drive_domclick(
|
||||||
recorder = _RunsRecorder()
|
recorder = _RunsRecorder()
|
||||||
db = MagicMock()
|
db = MagicMock()
|
||||||
lots = [MagicMock() for _ in range(lots_n)]
|
lots = [MagicMock() for _ in range(lots_n)]
|
||||||
|
|
||||||
|
async def _fetch_city(**kw: Any) -> list[Any]:
|
||||||
|
# fix/domclick-incremental-save: save_listings переехал внутрь цикла по корзинам и
|
||||||
|
# зовётся из колбэка on_bucket. Стаб, который просто возвращает лоты, не
|
||||||
|
# вызвал бы сохранение вовсе — фикстура проверяла бы мёртвую ветку.
|
||||||
|
# Одна корзина со всеми лотами: ровно один save_listings, как и было.
|
||||||
|
on_bucket = kw.get("on_bucket")
|
||||||
|
if on_bucket is not None and lots:
|
||||||
|
on_bucket(ROOM_BUCKETS[0], lots)
|
||||||
|
return lots
|
||||||
|
|
||||||
scraper = _ctx_scraper(
|
scraper = _ctx_scraper(
|
||||||
fetch_city=AsyncMock(return_value=lots),
|
fetch_city=_fetch_city,
|
||||||
blocked=blocked,
|
blocked=blocked,
|
||||||
geo_filtered=0,
|
geo_filtered=0,
|
||||||
fetch_errors=fetch_errors,
|
fetch_errors=fetch_errors,
|
||||||
|
|
|
||||||
|
|
@ -4868,6 +4868,7 @@ async def run_domclick_city_sweep(
|
||||||
region_code: int = DEFAULT_REGION_CODE,
|
region_code: int = DEFAULT_REGION_CODE,
|
||||||
resume_run_id: int | None = None,
|
resume_run_id: int | None = None,
|
||||||
cookies: dict[str, str] | None = None,
|
cookies: dict[str, str] | None = None,
|
||||||
|
watchdog_sec: int | None = None,
|
||||||
) -> DomClickCitySweepCounters:
|
) -> DomClickCitySweepCounters:
|
||||||
"""DomClick citywide sweep через BFF JSON API.
|
"""DomClick citywide sweep через BFF JSON API.
|
||||||
|
|
||||||
|
|
@ -4896,6 +4897,15 @@ async def run_domclick_city_sweep(
|
||||||
зависли на challenge). None (дефолт, сессии в БД нет/протухла) — прежнее
|
зависли на challenge). None (дефолт, сессии в БД нет/протухла) — прежнее
|
||||||
поведение, без инъекции.
|
поведение, без инъекции.
|
||||||
|
|
||||||
|
watchdog_sec (incremental-save): явный override расчётной формулы watchdog'а. Формула
|
||||||
|
ниже (buckets × pages × per_fetch + budget) не знает про бисекцию по цене
|
||||||
|
(ДомКлик режет offset на 2000, каждый лист пагинируется отдельно отдельным
|
||||||
|
деревом сплитов) и на больших городах (Москва ≈23 690 лотов против ЕКБ
|
||||||
|
≈6 300) занижена в разы — прод run 7344: за 2ч watchdog'а (11 100с) пройдено
|
||||||
|
2 бакета из 6. С инкрементальным сохранением (on_bucket, см. ниже) ранний
|
||||||
|
снос по watchdog больше не теряет собранное, поэтому вместо более точной
|
||||||
|
оценки — простой override: None (дефолт) — прежняя формула байт-в-байт.
|
||||||
|
|
||||||
Возвращает DomClickCitySweepCounters.
|
Возвращает DomClickCitySweepCounters.
|
||||||
"""
|
"""
|
||||||
# Гео-скоуп свипа больше не зашит в ЕКБ: его задаёт профиль региона (address
|
# Гео-скоуп свипа больше не зашит в ЕКБ: его задаёт профиль региона (address
|
||||||
|
|
@ -4909,12 +4919,18 @@ async def run_domclick_city_sweep(
|
||||||
_resolved_delay = request_delay_sec if request_delay_sec is not None else 6.0
|
_resolved_delay = request_delay_sec if request_delay_sec is not None else 6.0
|
||||||
counters = DomClickCitySweepCounters()
|
counters = DomClickCitySweepCounters()
|
||||||
|
|
||||||
# Watchdog: 6 buckets × pages × per_fetch + budget.
|
# Watchdog: 6 buckets × pages × per_fetch + budget. watchdog_sec (incremental-save) — явный
|
||||||
|
# override для больших городов, где формула занижена (см. докстринг выше).
|
||||||
_num_fetches = _DOMCLICK_NUM_BUCKETS * max(1, pages)
|
_num_fetches = _DOMCLICK_NUM_BUCKETS * max(1, pages)
|
||||||
_sweep_timeout = max(
|
if watchdog_sec is not None:
|
||||||
ANCHOR_TIMEOUT_SEC,
|
_sweep_timeout = int(watchdog_sec)
|
||||||
int(_num_fetches * (_resolved_delay + _DOMCLICK_PER_FETCH_S) + _DOMCLICK_SWEEP_BUDGET_S),
|
else:
|
||||||
)
|
_sweep_timeout = max(
|
||||||
|
ANCHOR_TIMEOUT_SEC,
|
||||||
|
int(
|
||||||
|
_num_fetches * (_resolved_delay + _DOMCLICK_PER_FETCH_S) + _DOMCLICK_SWEEP_BUDGET_S
|
||||||
|
),
|
||||||
|
)
|
||||||
|
|
||||||
# Мутируемый контейнер для захвата scraper-ссылки из замыкания.
|
# Мутируемый контейнер для захвата scraper-ссылки из замыкания.
|
||||||
_scraper_ref: list[DomClickScraper] = []
|
_scraper_ref: list[DomClickScraper] = []
|
||||||
|
|
@ -4991,14 +5007,68 @@ async def run_domclick_city_sweep(
|
||||||
_sweep_timeout,
|
_sweep_timeout,
|
||||||
)
|
)
|
||||||
|
|
||||||
lots: list[ScrapedLot] = []
|
# incremental-save: сохранение инкрементальное — save_listings зовётся из _on_bucket ПОСЛЕ
|
||||||
# #2406: собранное сохраняется ОДНИМ save_listings после всех корзин, поэтому
|
# КАЖДОГО room-бакета, а не одним save_listings в самом конце. Раньше снятие фазы
|
||||||
# корзина считается пройденной только если фаза дошла до конца (см. чекпоинт ниже).
|
# по watchdog'у (asyncio.wait_for TimeoutError) теряло ВСЁ собранное — на большой
|
||||||
_saved = False
|
# выдаче (Москва ≈23 690 лотов vs ЕКБ ≈6 300) формула watchdog'а систематически
|
||||||
|
# не укладывалась в отведённое время (прод run 7344: 2 бакета из 6 за 2ч).
|
||||||
|
_cancel_reason: str | None = None
|
||||||
|
|
||||||
|
def _on_bucket(bucket_key: str, bucket_lots: list[ScrapedLot]) -> None:
|
||||||
|
"""Инкрементальный save сразу после того, как room-бакет отфетчился.
|
||||||
|
|
||||||
|
fetch_city (serp.py) зовёт колбэк ВНЕ try/except конкретного бакета —
|
||||||
|
исключение отсюда прерывает обход ROOM_BUCKETS целиком (канал кооперативной
|
||||||
|
отмены/SIGTERM-дрейна), а не проглатывается generic except'ом бакета.
|
||||||
|
Sentinel-приём (RuntimeError("cancelled")/("shutdown")) — тот же, что в
|
||||||
|
run_cian_full_load._on_bucket (#1182 Phase 3a).
|
||||||
|
"""
|
||||||
|
nonlocal _checkpoint
|
||||||
|
if runs.is_cancelled(db, run_id):
|
||||||
|
logger.info(
|
||||||
|
"domclick-sweep run_id=%d: cancel detected in on_bucket (%s)",
|
||||||
|
run_id,
|
||||||
|
bucket_key,
|
||||||
|
)
|
||||||
|
raise RuntimeError("cancelled")
|
||||||
|
elif shutdown_requested():
|
||||||
|
logger.info(
|
||||||
|
"domclick-sweep run_id=%d: SIGTERM-drain — stopping at bucket %s",
|
||||||
|
run_id,
|
||||||
|
bucket_key,
|
||||||
|
)
|
||||||
|
raise RuntimeError("shutdown")
|
||||||
|
if bucket_lots:
|
||||||
|
# Имя города — из профиля региона, а не из сравнения с vestigial
|
||||||
|
# city_id: гео-скоп задаёт регион, он же знает, какой город штамповать.
|
||||||
|
# У области (50) city_name=None — одного города нет, угадывать нечего.
|
||||||
|
inserted, updated = save_listings(
|
||||||
|
db,
|
||||||
|
bucket_lots,
|
||||||
|
matcher=matcher,
|
||||||
|
region_code=region_code,
|
||||||
|
run_id=run_id,
|
||||||
|
city=_geo_profile.city_name,
|
||||||
|
)
|
||||||
|
counters.lots_fetched += len(bucket_lots)
|
||||||
|
counters.lots_inserted += inserted
|
||||||
|
counters.lots_updated += updated
|
||||||
|
# #3118, теперь на уровне бакета: done_buckets означает "собрано И
|
||||||
|
# сохранено" — чекпоинт пишем ТОЛЬКО пройдя cancel/shutdown-гейт выше и
|
||||||
|
# save_listings этого бакета, не из scraper.completed_buckets (см.
|
||||||
|
# комментарий у "Перенести счётчики" ниже — там раньше был баг #2 issue).
|
||||||
|
_checkpoint = sorted(set(_checkpoint) | {bucket_key})
|
||||||
|
runs.update_heartbeat(db, run_id, _payload())
|
||||||
|
logger.info(
|
||||||
|
"domclick-sweep run_id=%d: bucket %s saved lots=%d total=%d",
|
||||||
|
run_id,
|
||||||
|
bucket_key,
|
||||||
|
len(bucket_lots),
|
||||||
|
counters.lots_fetched,
|
||||||
|
)
|
||||||
|
|
||||||
async def _domclick_phase() -> None:
|
async def _domclick_phase() -> None:
|
||||||
"""Единственная citywide-фаза: fetch_city + save."""
|
"""Единственная citywide-фаза: fetch_city с инкрементальным save по бакетам."""
|
||||||
nonlocal lots, _saved
|
|
||||||
async with DomClickScraper(
|
async with DomClickScraper(
|
||||||
config,
|
config,
|
||||||
proxy_provider=proxy_provider,
|
proxy_provider=proxy_provider,
|
||||||
|
|
@ -5021,36 +5091,21 @@ async def run_domclick_city_sweep(
|
||||||
# источники и растёт неравномерно, — но за 30 суток каждая корзина
|
# источники и растёт неравномерно, — но за 30 суток каждая корзина
|
||||||
# получает порядка пяти стартов, чего достаточно для критерия приёмки
|
# получает порядка пяти стартов, чего достаточно для критерия приёмки
|
||||||
# «объявления с rooms >= 2 появились».
|
# «объявления с rooms >= 2 появились».
|
||||||
lots = await _scraper.fetch_city(
|
await _scraper.fetch_city(
|
||||||
city_id=city_id,
|
city_id=city_id,
|
||||||
rooms=rooms,
|
rooms=rooms,
|
||||||
pages=pages,
|
pages=pages,
|
||||||
start_bucket_index=run_id % len(ROOM_BUCKETS),
|
start_bucket_index=run_id % len(ROOM_BUCKETS),
|
||||||
skip_buckets=skip_buckets or None,
|
skip_buckets=skip_buckets or None,
|
||||||
|
on_bucket=_on_bucket,
|
||||||
)
|
)
|
||||||
counters.lots_fetched += len(lots)
|
|
||||||
if lots:
|
|
||||||
# Имя города — из профиля региона, а не из сравнения с vestigial
|
|
||||||
# city_id: гео-скоп задаёт регион, он же знает, какой город штамповать.
|
|
||||||
# У области (50) city_name=None — одного города нет, угадывать нечего.
|
|
||||||
_dc_city = _geo_profile.city_name
|
|
||||||
inserted, updated = save_listings(
|
|
||||||
db,
|
|
||||||
lots,
|
|
||||||
matcher=matcher,
|
|
||||||
region_code=region_code,
|
|
||||||
run_id=run_id,
|
|
||||||
city=_dc_city,
|
|
||||||
)
|
|
||||||
counters.lots_inserted += inserted
|
|
||||||
counters.lots_updated += updated
|
|
||||||
_saved = True
|
|
||||||
|
|
||||||
try:
|
try:
|
||||||
await asyncio.wait_for(_domclick_phase(), timeout=_sweep_timeout)
|
await asyncio.wait_for(_domclick_phase(), timeout=_sweep_timeout)
|
||||||
except TimeoutError:
|
except TimeoutError:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"domclick-sweep run_id=%d: SERP phase timed out after %ds — partial results",
|
"domclick-sweep run_id=%d: SERP phase timed out after %ds — partial results "
|
||||||
|
"(incremental-save: уже сохранены инкрементально, бакет за бакетом)",
|
||||||
run_id,
|
run_id,
|
||||||
_sweep_timeout,
|
_sweep_timeout,
|
||||||
)
|
)
|
||||||
|
|
@ -5080,6 +5135,16 @@ async def run_domclick_city_sweep(
|
||||||
ban_kind=ban_kind_of_exception(exc),
|
ban_kind=ban_kind_of_exception(exc),
|
||||||
)
|
)
|
||||||
return counters
|
return counters
|
||||||
|
except RuntimeError as exc:
|
||||||
|
# on_bucket кидает RuntimeError("cancelled") при кооперативной отмене,
|
||||||
|
# RuntimeError("shutdown") при SIGTERM-дрейне — тот же sentinel-приём, что в
|
||||||
|
# run_cian_full_load._on_bucket (#1182 Phase 3a). Прочие RuntimeError —
|
||||||
|
# обычная поломка фазы, ведём себя как под generic except ниже.
|
||||||
|
if str(exc) in ("cancelled", "shutdown"):
|
||||||
|
_cancel_reason = str(exc)
|
||||||
|
else:
|
||||||
|
logger.exception("domclick-sweep run_id=%d: SERP phase failed", run_id)
|
||||||
|
counters.errors_count += 1
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.exception("domclick-sweep run_id=%d: SERP phase failed", run_id)
|
logger.exception("domclick-sweep run_id=%d: SERP phase failed", run_id)
|
||||||
counters.errors_count += 1
|
counters.errors_count += 1
|
||||||
|
|
@ -5095,15 +5160,13 @@ async def run_domclick_city_sweep(
|
||||||
# скрейпер живая, а его счётчик показывает, докуда прогон дошёл.
|
# скрейпер живая, а его счётчик показывает, докуда прогон дошёл.
|
||||||
counters.buckets_completed = _s.buckets_completed
|
counters.buckets_completed = _s.buckets_completed
|
||||||
counters.buckets_total = _s.buckets_total
|
counters.buckets_total = _s.buckets_total
|
||||||
# #3118: чекпоинт = унаследованное ∪ завершённое в этом прогоне.
|
# incremental-save: _checkpoint СЮДА больше не мержится из _s.completed_buckets — он
|
||||||
# Пишем heartbeat'ом СЕЙЧАС (мерж jsonb) — финализаторы ключ не
|
# пишется инкрементально внутри _on_bucket, СРАЗУ после save_listings этого
|
||||||
# затирают, и оборванный болезнью финализации прогон его не теряет.
|
# бакета. _s.completed_buckets на уровне скрейпера отмечает бакет как "фетч
|
||||||
# #2406: только если лоты сохранены. Снятая watchdog'ом или упавшая фаза
|
# прошёл" ДО вызова on_bucket (см. serp.py): если on_bucket поймал
|
||||||
# не дошла до save_listings — у её корзин в БД ноль строк, а отметка
|
# cancel/shutdown ДО save для этого самого бакета, _s.completed_buckets
|
||||||
# «пройдена» заставила бы следующий прогон пропустить их навсегда
|
# включил бы его, а _checkpoint — честно нет (лоты не сохранены). Мердж
|
||||||
# (механизм разобран в миграции 308, из-за него выключены свипы 77/50).
|
# отсюда воспроизвёл бы старый баг — чекпоинт врёт про несохранённые бакеты.
|
||||||
if _saved:
|
|
||||||
_checkpoint = sorted(skip_buckets | set(_s.completed_buckets))
|
|
||||||
runs.update_heartbeat(db, run_id, _payload())
|
runs.update_heartbeat(db, run_id, _payload())
|
||||||
counters.bucket_start_index = _s.bucket_start_index
|
counters.bucket_start_index = _s.bucket_start_index
|
||||||
|
|
||||||
|
|
@ -5111,6 +5174,32 @@ async def run_domclick_city_sweep(
|
||||||
counters.pages_fetched = _num_fetches
|
counters.pages_fetched = _num_fetches
|
||||||
runs.update_heartbeat(db, run_id, _payload())
|
runs.update_heartbeat(db, run_id, _payload())
|
||||||
|
|
||||||
|
# incremental-save: кооперативная отмена/SIGTERM-дрейн, пойманные в on_bucket — партиал уже
|
||||||
|
# сохранён инкрементально, финализируем как run_cian_full_load (mark_done partial,
|
||||||
|
# не honest-status ниже: обрыв тут known-signal, а не "прогон не доделал сам").
|
||||||
|
if _cancel_reason == "cancelled":
|
||||||
|
logger.info(
|
||||||
|
"domclick-sweep run_id=%d: cancelled — partial results lots=%d (ins=%d/upd=%d)",
|
||||||
|
run_id,
|
||||||
|
counters.lots_fetched,
|
||||||
|
counters.lots_inserted,
|
||||||
|
counters.lots_updated,
|
||||||
|
)
|
||||||
|
runs.mark_done(db, run_id, _payload())
|
||||||
|
return counters
|
||||||
|
if _cancel_reason == "shutdown":
|
||||||
|
logger.info(
|
||||||
|
"domclick-sweep run_id=%d: SIGTERM-drain — partial results lots=%d (ins=%d/upd=%d)",
|
||||||
|
run_id,
|
||||||
|
counters.lots_fetched,
|
||||||
|
counters.lots_inserted,
|
||||||
|
counters.lots_updated,
|
||||||
|
)
|
||||||
|
_drain = {**_payload(), "interrupted": 1}
|
||||||
|
runs.update_heartbeat(db, run_id, _drain)
|
||||||
|
runs.mark_done(db, run_id, _drain)
|
||||||
|
return counters
|
||||||
|
|
||||||
# ── ЧЕСТНЫЙ СТАТУС (#1968, ужесточён #2657) ───────────────────────────
|
# ── ЧЕСТНЫЙ СТАТУС (#1968, ужесточён #2657) ───────────────────────────
|
||||||
# Распознанный QRATOR-блок НИКОГДА не даёт done. Домклик тут структурно
|
# Распознанный QRATOR-блок НИКОГДА не даёт done. Домклик тут структурно
|
||||||
# отличается от cian/yandex (#2625/#2642): там независимые anchor'ы, и
|
# отличается от cian/yandex (#2625/#2642): там независимые anchor'ы, и
|
||||||
|
|
|
||||||
|
|
@ -1270,6 +1270,9 @@ async def _job_domclick_city_sweep(
|
||||||
request_delay_sec=float(params.get("request_delay_sec", 6.0)),
|
request_delay_sec=float(params.get("request_delay_sec", 6.0)),
|
||||||
region_code=_resolve_region_code(params),
|
region_code=_resolve_region_code(params),
|
||||||
resume_run_id=_pick_resume(db, run_id),
|
resume_run_id=_pick_resume(db, run_id),
|
||||||
|
watchdog_sec=(
|
||||||
|
int(params["watchdog_sec"]) if params.get("watchdog_sec") is not None else None
|
||||||
|
),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -459,6 +459,7 @@ class DomClickScraper(BaseScraper):
|
||||||
pages: int = 100,
|
pages: int = 100,
|
||||||
start_bucket_index: int = 0,
|
start_bucket_index: int = 0,
|
||||||
skip_buckets: set[str] | None = None,
|
skip_buckets: set[str] | None = None,
|
||||||
|
on_bucket: Callable[[str, list[ScrapedLot]], None] | None = None,
|
||||||
) -> list[ScrapedLot]:
|
) -> list[ScrapedLot]:
|
||||||
"""Citywide sweep через BFF JSON API.
|
"""Citywide sweep через BFF JSON API.
|
||||||
|
|
||||||
|
|
@ -492,6 +493,17 @@ class DomClickScraper(BaseScraper):
|
||||||
buckets_total при этом = числу корзин В ЭТОМ прогоне (без
|
buckets_total при этом = числу корзин В ЭТОМ прогоне (без
|
||||||
скипнутых) — иначе honest-status читал бы возобновлённый прогон
|
скипнутых) — иначе honest-status читал бы возобновлённый прогон
|
||||||
как вечно-частичный.
|
как вечно-частичный.
|
||||||
|
on_bucket: колбэк инкрементального сохранения — большая выдача теряла
|
||||||
|
ВСЁ собранное при снятии по watchdog, единственный save был в самом
|
||||||
|
конце (fix/domclick-incremental-save). Зовётся СИНХРОННО сразу после того, как
|
||||||
|
бакет отработал успешно (после DomClickBlockedError/generic
|
||||||
|
Exception — НЕ зовётся), аргументы: имя бакета + лоты ИМЕННО
|
||||||
|
этого бакета (не накопленный out_lots). Вызов стоит ВНЕ
|
||||||
|
try/except этого бакета: исключение из колбэка (кооперативная
|
||||||
|
отмена/SIGTERM-дрейн — см. run_domclick_city_sweep) обязано
|
||||||
|
прервать цикл по ROOM_BUCKETS, а не быть проглоченным generic
|
||||||
|
except'ом. on_bucket=None (дефолт) — поведение прежнее,
|
||||||
|
байт-в-байт.
|
||||||
|
|
||||||
Returns:
|
Returns:
|
||||||
Дедуплицированный по source_id список ScrapedLot.
|
Дедуплицированный по source_id список ScrapedLot.
|
||||||
|
|
@ -548,6 +560,7 @@ class DomClickScraper(BaseScraper):
|
||||||
city_id,
|
city_id,
|
||||||
pages,
|
pages,
|
||||||
)
|
)
|
||||||
|
_bucket_start_len = len(out_lots)
|
||||||
try:
|
try:
|
||||||
await self._sweep_bucket(
|
await self._sweep_bucket(
|
||||||
fetcher=fetcher,
|
fetcher=fetcher,
|
||||||
|
|
@ -591,6 +604,11 @@ class DomClickScraper(BaseScraper):
|
||||||
continue
|
continue
|
||||||
self.buckets_completed += 1
|
self.buckets_completed += 1
|
||||||
self.completed_buckets.append(bucket)
|
self.completed_buckets.append(bucket)
|
||||||
|
# incremental-save: колбэк ВНЕ try/except этого бакета — исключение (канал
|
||||||
|
# кооперативной отмены/SIGTERM-дрейна в pipeline.run_domclick_city_sweep)
|
||||||
|
# обязано прервать цикл, а не попасть в generic except выше.
|
||||||
|
if on_bucket is not None:
|
||||||
|
on_bucket(bucket, out_lots[_bucket_start_len:])
|
||||||
|
|
||||||
logger.info(
|
logger.info(
|
||||||
"domklik: fetch_city done city_id=%d total=%d buckets=%d/%d "
|
"domklik: fetch_city done city_id=%d total=%d buckets=%d/%d "
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue