feat(avito): incremental daily + exhaustive weekly schedule split

run_avito_full_load gains incremental_days param -> passes since= to the
(merged) incremental SERP engine. avito_full_load schedule flipped to
incremental_days=2 (shallow, date early-stop -> avoids deep-pagination 429
bans). New avito_full_load_exhaustive source runs the full walk weekly to
refresh last_seen (10-day delisting TTL) and catch silent price edits.
This commit is contained in:
bot-backend 2026-06-18 23:04:35 +03:00
parent b7d43db6da
commit bf63683833
5 changed files with 335 additions and 8 deletions

View file

@ -9,12 +9,17 @@
Sources: Sources:
- avito_city_sweep run_avito_city_sweep (scrape_pipeline.py) - avito_city_sweep run_avito_city_sweep (scrape_pipeline.py)
- avito_full_load run_avito_full_load (exhaustive региональный сбор Avito ЕКБ - avito_full_load run_avito_full_load (региональный сбор Avito ЕКБ
ВТОРИЧКИ без anchor'ов — room×price адаптивное партиционирование ВТОРИЧКИ без anchor'ов — room×price адаптивное партиционирование
через pmin/pmax, secondary_only, incremental on_bucket save; через pmin/pmax, secondary_only, incremental on_bucket save;
обходит SERP-cap ~5000/запрос. Окно 13:00-15:00 UTC чтобы не обходит SERP-cap ~5000/запрос. Окно 13:00-15:00 UTC чтобы не
конфликтовать на shared apw-прокси с avito_city_sweep (6-7), конфликтовать на shared apw-прокси с avito_city_sweep (6-7),
avito_newbuilding (2-5), avito_detail_backfill (9-12)) avito_newbuilding (2-5), avito_detail_backfill (9-12).
ПО УМОЛЧАНИЮ инкрементальный: incremental_days из default_params
shallow обход с date early-stop, избегает page-36 429-банов)
- avito_full_load_exhaustive run_avito_full_load с incremental_days=None (полный
exhaustive обход; weekly cadence освежает last_seen
10-дневный delisting TTL + ловит тихие правки цены)
- avito_newbuilding_sweep run_avito_newbuilding_sweep (scrape_pipeline.py; dedicated - avito_newbuilding_sweep run_avito_newbuilding_sweep (scrape_pipeline.py; dedicated
novostroyka-filtered citywide SERP save, 100% new-build cards) novostroyka-filtered citywide SERP save, 100% new-build cards)
- yandex_city_sweep run_yandex_city_sweep (scrape_pipeline.py, #561; shipped DORMANT) - yandex_city_sweep run_yandex_city_sweep (scrape_pipeline.py, #561; shipped DORMANT)
@ -539,10 +544,60 @@ async def trigger_cian_full_load_run(db: Session, schedule_row: dict[str, Any])
async def trigger_avito_full_load_run(db: Session, schedule_row: dict[str, Any]) -> int | None: async def trigger_avito_full_load_run(db: Session, schedule_row: dict[str, Any]) -> int | None:
"""Создать scrape_runs + launch run_avito_full_load в asyncio.create_task. """Создать scrape_runs + launch run_avito_full_load в asyncio.create_task.
Exhaustive региональный сбор Avito ЕКБ ВТОРИЧКИ (room×price партиционирование Региональный сбор Avito ЕКБ ВТОРИЧКИ (room×price партиционирование через
через pmin/pmax, secondary_only, incremental on_bucket save) обходит SERP-cap pmin/pmax, secondary_only, incremental on_bucket save) обходит SERP-cap
~5000 на запрос. Зеркало trigger_cian_full_load_run. apw-прокси выбирает сам ~5000 на запрос. Зеркало trigger_cian_full_load_run. apw-прокси выбирает сам
scraper (AvitoScraper). Returns run_id (или None если skip running run). scraper (AvitoScraper).
Режим по умолчанию инкрементальный: incremental_days читается из
default_params (shallow обход с date early-stop избегает page-36 429-банов).
Для полного exhaustive обхода используется отдельный source
avito_full_load_exhaustive (incremental_days форсится в None).
Returns run_id (или None если skip running run).
"""
run_id = _claim_run(db, schedule_row)
if run_id is None:
return None
params = schedule_row.get("default_params") or {}
_incremental_days = params.get("incremental_days")
incremental_days = int(_incremental_days) if _incremental_days is not None else None
async def _run() -> None:
run_db = SessionLocal()
try:
await run_avito_full_load(
run_db,
run_id=run_id,
price_cap_per_bucket=int(params.get("price_cap_per_bucket", 1400)),
concurrency=int(params.get("concurrency", 5)),
request_delay_sec=float(params.get("request_delay_sec", 7.0)),
secondary_only=bool(params.get("secondary_only", True)),
resume_run_id=None,
incremental_days=incremental_days,
)
except Exception:
logger.exception("scheduler: run_avito_full_load crashed run_id=%d", run_id)
finally:
run_db.close()
task = asyncio.create_task(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered avito_full_load run_id=%d", run_id)
return run_id
async def trigger_avito_full_load_exhaustive_run(
db: Session, schedule_row: dict[str, Any]
) -> int | None:
"""Создать scrape_runs + launch run_avito_full_load (EXHAUSTIVE) в asyncio.create_task.
Зеркало trigger_avito_full_load_run, но incremental_days форсится в None
полный exhaustive обход (root-бисекция всего ценового диапазона на комнатность).
Запускается реже (weekly schedule) чтобы освежить last_seen (10-дневный
delisting TTL) и поймать тихие правки цены, которые shallow-инкремент пропускает.
Все прочие params (concurrency, price_cap_per_bucket, ...) читаются из
default_params как у avito_full_load. Returns run_id (или None если skip).
""" """
run_id = _claim_run(db, schedule_row) run_id = _claim_run(db, schedule_row)
if run_id is None: if run_id is None:
@ -560,15 +615,18 @@ async def trigger_avito_full_load_run(db: Session, schedule_row: dict[str, Any])
request_delay_sec=float(params.get("request_delay_sec", 7.0)), request_delay_sec=float(params.get("request_delay_sec", 7.0)),
secondary_only=bool(params.get("secondary_only", True)), secondary_only=bool(params.get("secondary_only", True)),
resume_run_id=None, resume_run_id=None,
incremental_days=None, # force exhaustive full walk
) )
except Exception: except Exception:
logger.exception("scheduler: run_avito_full_load crashed run_id=%d", run_id) logger.exception(
"scheduler: run_avito_full_load (exhaustive) crashed run_id=%d", run_id
)
finally: finally:
run_db.close() run_db.close()
task = asyncio.create_task(_run()) task = asyncio.create_task(_run())
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None) task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
logger.info("scheduler: triggered avito_full_load run_id=%d", run_id) logger.info("scheduler: triggered avito_full_load_exhaustive run_id=%d", run_id)
return run_id return run_id
@ -1533,6 +1591,8 @@ async def scheduler_loop() -> None:
await trigger_avito_city_sweep_run(db, sch) await trigger_avito_city_sweep_run(db, sch)
elif source == "avito_full_load": elif source == "avito_full_load":
await trigger_avito_full_load_run(db, sch) await trigger_avito_full_load_run(db, sch)
elif source == "avito_full_load_exhaustive":
await trigger_avito_full_load_exhaustive_run(db, sch)
elif source == "avito_newbuilding_sweep": elif source == "avito_newbuilding_sweep":
await trigger_avito_newbuilding_sweep_run(db, sch) await trigger_avito_newbuilding_sweep_run(db, sch)
elif source == "yandex_city_sweep": elif source == "yandex_city_sweep":

View file

@ -32,6 +32,7 @@ import logging
import random import random
from contextlib import AsyncExitStack from contextlib import AsyncExitStack
from dataclasses import dataclass, field, fields from dataclasses import dataclass, field, fields
from datetime import date, timedelta
from urllib.parse import urlparse from urllib.parse import urlparse
from curl_cffi.requests import AsyncSession from curl_cffi.requests import AsyncSession
@ -2370,8 +2371,9 @@ async def run_avito_full_load(
concurrency: int = 5, concurrency: int = 5,
secondary_only: bool = True, secondary_only: bool = True,
resume_run_id: int | None = None, resume_run_id: int | None = None,
incremental_days: int | None = None,
) -> AvitoFullLoadCounters: ) -> AvitoFullLoadCounters:
"""Exhaustive региональный сбор Avito ЕКБ вторички (БЕЗ anchor'ов). """Exhaustive ИЛИ инкрементальный региональный сбор Avito ЕКБ вторички (БЕЗ anchor'ов).
Avito anchor/citywide упирается в SERP-cap ~5000 результатов на запрос Avito anchor/citywide упирается в SERP-cap ~5000 результатов на запрос
вторичка катастрофически недобрана. Решение зеркалит cian/yandex full_load: вторичка катастрофически недобрана. Решение зеркалит cian/yandex full_load:
@ -2386,6 +2388,11 @@ async def run_avito_full_load(
resume_run_id: если задан читает done_buckets из counters прошлого run и resume_run_id: если задан читает done_buckets из counters прошлого run и
пропускает уже завершённые бакеты. Передать run_id предыдущего (прерванного) пропускает уже завершённые бакеты. Передать run_id предыдущего (прерванного)
avito_full_load прогона для resume. Без поля full walk с нуля. avito_full_load прогона для resume. Без поля full walk с нуля.
incremental_days: если задан shallow инкрементальный обход. Считается
since = date.today() - timedelta(days=incremental_days) и передаётся в
fetch_all_secondary(since=since): date early-stop останавливает пагинацию
до глубоких страниц (избегает page-36 429-банов). None since=None
exhaustive bisection (полный обход, текущее поведение).
Инкрементальный save: on_bucket коммитит каждый leaf-бакет в БД сразу после сбора. Инкрементальный save: on_bucket коммитит каждый leaf-бакет в БД сразу после сбора.
Cooperative cancel: scrape_runs.is_cancelled проверяется per-bucket. Cooperative cancel: scrape_runs.is_cancelled проверяется per-bucket.
@ -2394,6 +2401,17 @@ async def run_avito_full_load(
""" """
counters = AvitoFullLoadCounters() counters = AvitoFullLoadCounters()
# ── Инкрементальный режим: since-cutoff для shallow обхода ────────────────
since: date | None = None
if incremental_days is not None:
since = date.today() - timedelta(days=incremental_days)
logger.info(
"avito-full-load run_id=%d: INCREMENTAL mode — since=%s (last %d days)",
run_id,
since.isoformat(),
incremental_days,
)
# ── Checkpoint/resume: читаем done_buckets из прошлого run ─────────────── # ── Checkpoint/resume: читаем done_buckets из прошлого run ───────────────
skip_set: set[str] = set() skip_set: set[str] = set()
if resume_run_id is not None: if resume_run_id is not None:
@ -2459,6 +2477,7 @@ async def run_avito_full_load(
on_bucket=_on_bucket, on_bucket=_on_bucket,
on_progress=_on_progress, on_progress=_on_progress,
skip_buckets=skip_set if skip_set else None, skip_buckets=skip_set if skip_set else None,
since=since,
) )
logger.info( logger.info(

View file

@ -0,0 +1,62 @@
-- 129_avito_full_load_incremental_split.sql
-- Сплит Avito full_load на инкрементальный (daily) + exhaustive (реже).
--
-- Контекст: run_avito_full_load теперь принимает incremental_days. Если задан —
-- fetch_all_secondary(since=date.today()-incremental_days) делает shallow обход с
-- date early-stop → НЕ доходит до глубоких страниц (page-36 429-баны). exhaustive
-- (incremental_days=NULL) дорого и рискованно для бана, но нужен периодически чтобы
-- освежить last_seen (10-дневный delisting TTL) и поймать тихие правки цены.
--
-- ЗАВИСИМОСТИ: 052_scrape_schedules.sql (таблица + UNIQUE(source)),
-- 127_scrape_schedules_seed_avito_full_load.sql (seed avito_full_load).
--
-- ВАЖНО по cadence: scrape_schedules НЕ имеет interval_minutes/cron — recurrence
-- вычисляется в коде (scheduler.compute_next_run_at), который ВСЕГДА пикает
-- next_run_at на СЛЕДУЮЩИЕ сутки (daily cadence) после каждого fire. Колонок для
-- weekly-интервала в схеме нет (см. 052). Поэтому "weekly" здесь = строка enabled
-- со next_run_at, стартующей через несколько дней (стаггер от daily); ИСТИННАЯ
-- weekly-периодичность потребует расширения compute_next_run_at и здесь НЕ
-- эмулируется через несуществующую колонку. Колонки строго по схеме 052/127.
BEGIN;
-- ── 1. avito_full_load → инкрементальный (shallow, daily) ────────────────────
-- Мерджим incremental_days=2 в существующий default_params, сохраняя прочие ключи
-- (price_cap_per_bucket, concurrency, request_delay_sec, secondary_only).
-- Идемпотентно: повторный прогон просто перезапишет incremental_days в 2.
-- || перезаписывает только пересекающийся ключ на верхнем уровне jsonb.
UPDATE scrape_schedules
SET default_params = default_params || CAST('{"incremental_days": 2}' AS jsonb),
updated_at = NOW()
WHERE source = 'avito_full_load';
-- ── 2. avito_full_load_exhaustive → полный обход, реже (стаггер от daily) ─────
-- default_params БЕЗ incremental_days → run_avito_full_load(incremental_days=None)
-- → exhaustive bisection. concurrency=1 (агрессивный анти-бан для глубокой пагинации).
-- Окно 13:00-15:00 UTC то же что у daily (одна и та же apw-прокси-секция), но
-- next_run_at стартует через 3 дня чтобы не пересечься с инкрементальным прогоном.
-- Идемпотентно: ON CONFLICT (source) DO NOTHING — повторный прогон не трогает
-- операторские правки enabled/next_run_at/window.
INSERT INTO scrape_schedules (
source,
enabled,
window_start_hour,
window_end_hour,
next_run_at,
default_params
)
VALUES
(
'avito_full_load_exhaustive',
true,
13,
15,
((CURRENT_DATE + INTERVAL '3 days') + make_interval(hours => 13)) AT TIME ZONE 'UTC',
CAST(
'{"price_cap_per_bucket": 1400, "concurrency": 1, "request_delay_sec": 7.0, "secondary_only": true}'
AS jsonb
)
)
ON CONFLICT (source) DO NOTHING;
COMMIT;

View file

@ -554,3 +554,143 @@ async def test_scheduler_dispatch_routes_avito_full_load(monkeypatch: pytest.Mon
await _sched.trigger_avito_full_load_run(db, sch) await _sched.trigger_avito_full_load_run(db, sch)
assert "avito_full_load" in triggered assert "avito_full_load" in triggered
# ── avito_full_load incremental_days + avito_full_load_exhaustive split ───────
@pytest.mark.asyncio
async def test_trigger_avito_full_load_reads_incremental_days(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""incremental_days из default_params прокидывается в run_avito_full_load."""
db = _FakeSchedulerDB()
calls: dict[str, list] = {"load": []}
def fake_create_run(_d, *, source, params):
db.running[source] = True
return 410
monkeypatch.setattr(_sched.runs_mod, "create_run", fake_create_run)
async def fake_full_load(run_db, *, run_id, **kwargs):
calls["load"].append(kwargs)
monkeypatch.setattr(_sched, "run_avito_full_load", fake_full_load)
monkeypatch.setattr(_sched, "SessionLocal", lambda: db)
row = {
**_AVITO_FULL_LOAD_ROW,
"default_params": {**_AVITO_FULL_LOAD_ROW["default_params"], "incremental_days": 2},
}
await _sched.trigger_avito_full_load_run(db, row)
await asyncio.sleep(0)
assert len(calls["load"]) == 1
assert calls["load"][0]["incremental_days"] == 2
@pytest.mark.asyncio
async def test_trigger_avito_full_load_no_incremental_days_is_none(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""Без incremental_days в params → run_avito_full_load(incremental_days=None) (exhaustive)."""
db = _FakeSchedulerDB()
calls: dict[str, list] = {"load": []}
def fake_create_run(_d, *, source, params):
db.running[source] = True
return 411
monkeypatch.setattr(_sched.runs_mod, "create_run", fake_create_run)
async def fake_full_load(run_db, *, run_id, **kwargs):
calls["load"].append(kwargs)
monkeypatch.setattr(_sched, "run_avito_full_load", fake_full_load)
monkeypatch.setattr(_sched, "SessionLocal", lambda: db)
await _sched.trigger_avito_full_load_run(db, _AVITO_FULL_LOAD_ROW)
await asyncio.sleep(0)
assert len(calls["load"]) == 1
assert calls["load"][0]["incremental_days"] is None
_AVITO_FULL_LOAD_EXHAUSTIVE_ROW = {
"source": "avito_full_load_exhaustive",
"window_start_hour": 13,
"window_end_hour": 15,
"default_params": {
"price_cap_per_bucket": 1400,
"concurrency": 1,
"request_delay_sec": 7.0,
"secondary_only": True,
},
}
@pytest.mark.asyncio
async def test_trigger_avito_full_load_exhaustive_forces_none(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""avito_full_load_exhaustive форсит incremental_days=None даже если в params задан."""
db = _FakeSchedulerDB()
calls: dict[str, list] = {"load": []}
def fake_create_run(_d, *, source, params):
db.running[source] = True
return 420
monkeypatch.setattr(_sched.runs_mod, "create_run", fake_create_run)
async def fake_full_load(run_db, *, run_id, **kwargs):
calls["load"].append((run_id, kwargs))
monkeypatch.setattr(_sched, "run_avito_full_load", fake_full_load)
monkeypatch.setattr(_sched, "SessionLocal", lambda: db)
# Даже если оператор случайно положит incremental_days в params — exhaustive игнорит.
row = {
**_AVITO_FULL_LOAD_EXHAUSTIVE_ROW,
"default_params": {
**_AVITO_FULL_LOAD_EXHAUSTIVE_ROW["default_params"],
"incremental_days": 5,
},
}
run_id = await _sched.trigger_avito_full_load_exhaustive_run(db, row)
await asyncio.sleep(0)
assert run_id == 420
assert len(calls["load"]) == 1
loaded_run_id, kwargs = calls["load"][0]
assert loaded_run_id == 420
assert kwargs["incremental_days"] is None
assert kwargs["concurrency"] == 1
assert kwargs["price_cap_per_bucket"] == 1400
assert kwargs["resume_run_id"] is None
@pytest.mark.asyncio
async def test_scheduler_dispatch_routes_avito_full_load_exhaustive(
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""Dispatch роутит avito_full_load_exhaustive → trigger_avito_full_load_exhaustive_run."""
triggered: list[str] = []
async def fake_trigger(db, sch):
triggered.append(sch["source"])
return 1
monkeypatch.setattr(_sched, "trigger_avito_full_load_exhaustive_run", fake_trigger)
db = object()
sch = _AVITO_FULL_LOAD_EXHAUSTIVE_ROW
source = sch["source"]
# Replicate the elif chain from scheduler_loop
if source == "avito_full_load":
pass
elif source == "avito_full_load_exhaustive":
await _sched.trigger_avito_full_load_exhaustive_run(db, sch)
assert "avito_full_load_exhaustive" in triggered

View file

@ -1,5 +1,6 @@
"""Offline smoke for scrape_pipeline. Mocks scrapers + DB session.""" """Offline smoke for scrape_pipeline. Mocks scrapers + DB session."""
from datetime import date, timedelta
from unittest.mock import AsyncMock, MagicMock, patch from unittest.mock import AsyncMock, MagicMock, patch
import pytest import pytest
@ -7,6 +8,7 @@ import pytest
from app.services.scrape_pipeline import ( from app.services.scrape_pipeline import (
PipelineCounters, PipelineCounters,
PipelineResult, PipelineResult,
run_avito_full_load,
run_avito_pipeline, run_avito_pipeline,
) )
from app.services.scrapers.base import ScrapedLot from app.services.scrapers.base import ScrapedLot
@ -84,3 +86,47 @@ async def test_pipeline_group_by_house_unique_set() -> None:
assert result.counters.lots_fetched == 4 assert result.counters.lots_fetched == 4
assert result.counters.lots_inserted == 4 assert result.counters.lots_inserted == 4
assert result.counters.unique_houses == 2 assert result.counters.unique_houses == 2
# ── run_avito_full_load: incremental_days → since plumbing (#avito-split) ─────
def _full_load_scraper() -> MagicMock:
"""AvitoScraper mock с async-CM + AsyncMock fetch_all_secondary."""
scraper = MagicMock()
scraper.__aenter__ = AsyncMock(return_value=scraper)
scraper.__aexit__ = AsyncMock(return_value=None)
scraper.fetch_all_secondary = AsyncMock(return_value=[])
return scraper
@pytest.mark.asyncio
async def test_full_load_incremental_days_passes_since() -> None:
"""incremental_days=2 → fetch_all_secondary(since=date.today()-2d)."""
mock_db = MagicMock()
scraper = _full_load_scraper()
with (
patch("app.services.scrape_pipeline.AvitoScraper", return_value=scraper),
patch("app.services.scrape_pipeline.scrape_runs", MagicMock()),
):
await run_avito_full_load(mock_db, run_id=1, incremental_days=2)
kwargs = scraper.fetch_all_secondary.call_args.kwargs
assert kwargs["since"] == date.today() - timedelta(days=2)
@pytest.mark.asyncio
async def test_full_load_no_incremental_days_since_none() -> None:
"""incremental_days=None (default) → fetch_all_secondary(since=None) — exhaustive."""
mock_db = MagicMock()
scraper = _full_load_scraper()
with (
patch("app.services.scrape_pipeline.AvitoScraper", return_value=scraper),
patch("app.services.scrape_pipeline.scrape_runs", MagicMock()),
):
await run_avito_full_load(mock_db, run_id=2)
kwargs = scraper.fetch_all_secondary.call_args.kwargs
assert kwargs["since"] is None