fix(scrapers): cian full_load anti-zombie — per-fetch wait_for + heartbeat внутри бакета (#1949)
Два связанных фикса из диагностики 2026-06-27: ФИКС C (run_cian_full_load, per-fetch timeout): monkey-patch _browser.fetch → asyncio.wait_for(_f(url), timeout=cian_full_load_per_fetch_timeout_s=90s). BrowserFetcher.fetch имеет httpx-timeout 120s, но browser-сервис иногда виснет так что httpx не получает ответ (нет EOF) → fetch висит → asyncio.gather блокирует весь bucket → heartbeat_at не обновляется → reap_zombies убивает живой run. wait_for(90s) отменяет зависший fetch → TimeoutError → _fetch_page_html ловит как Exception → None → _one_page → [] → gather завершается нормально. За флагом (cian_full_load_per_fetch_timeout_s=0 = отключить, >0 = включить). ФИКС D (background heartbeat): asyncio.create_task(_background_heartbeat()) обновляет heartbeat_at каждые 60s независимо от on_bucket/on_progress прогресса. Без него пустые bucket или зависший gather замораживает heartbeat на несколько часов. Task отменяется в finally → не течёт при return/raise/cancel. Новая настройка (ENV): CIAN_FULL_LOAD_PER_FETCH_TIMEOUT_S=90.0 — timeout одного browser-fetch
This commit is contained in:
parent
757f6486f5
commit
e24d3cf4be
3 changed files with 290 additions and 0 deletions
|
|
@ -384,6 +384,13 @@ class Settings(BaseSettings):
|
||||||
# ENV: CIAN_PROXY_MAX_ROTATIONS.
|
# ENV: CIAN_PROXY_MAX_ROTATIONS.
|
||||||
cian_proxy_max_rotations: int = 4
|
cian_proxy_max_rotations: int = 4
|
||||||
|
|
||||||
|
# #1949: per-fetch hard-cancel для browser-fetch'ей внутри bucket-gather Cian full_load.
|
||||||
|
# asyncio.wait_for(fetch, timeout=X) отменяет зависший fetch → TimeoutError →
|
||||||
|
# _fetch_page_html ловит как Exception → None → _one_page → [] → gather завершается.
|
||||||
|
# 0 = отключить (старое поведение). Дефолт 90.0s щедрее nормального fetch ~12s.
|
||||||
|
# ENV: CIAN_FULL_LOAD_PER_FETCH_TIMEOUT_S.
|
||||||
|
cian_full_load_per_fetch_timeout_s: float = 90.0
|
||||||
|
|
||||||
@property
|
@property
|
||||||
def cian_proxy_url(self) -> str | None:
|
def cian_proxy_url(self) -> str | None:
|
||||||
"""Прокси для Cian-скраперов. CIAN_PROXY_URL > scraper_proxy_url (fallback)."""
|
"""Прокси для Cian-скраперов. CIAN_PROXY_URL > scraper_proxy_url (fallback)."""
|
||||||
|
|
|
||||||
|
|
@ -33,6 +33,7 @@ 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 datetime import date, timedelta
|
||||||
|
from typing import Any
|
||||||
from urllib.parse import urlparse
|
from urllib.parse import urlparse
|
||||||
|
|
||||||
from curl_cffi.requests import AsyncSession
|
from curl_cffi.requests import AsyncSession
|
||||||
|
|
@ -2535,10 +2536,61 @@ async def run_cian_full_load(
|
||||||
"""Heartbeat per room-bucket (на случай пустых бакетов без on_bucket вызовов)."""
|
"""Heartbeat per room-bucket (на случай пустых бакетов без on_bucket вызовов)."""
|
||||||
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
|
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
|
||||||
|
|
||||||
|
# #1949: фоновый heartbeat — обновляет heartbeat_at каждые 60s независимо от on_bucket.
|
||||||
|
# Без него зависший asyncio.gather внутри bucket-loop замораживает heartbeat на часы
|
||||||
|
# → reap_zombies помечает живой run zombie. Отменяется в finally при любом завершении.
|
||||||
|
async def _background_heartbeat() -> None:
|
||||||
|
"""Периодический heartbeat каждые 60s, не завязанный на прогресс бакетов."""
|
||||||
|
while True:
|
||||||
|
await asyncio.sleep(60.0)
|
||||||
|
try:
|
||||||
|
scrape_runs.update_heartbeat(db, run_id, counters.to_dict())
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
|
||||||
|
_hb_task: asyncio.Task[None] | None = None
|
||||||
try:
|
try:
|
||||||
|
_hb_task = asyncio.create_task(_background_heartbeat())
|
||||||
|
|
||||||
async with CianScraper() as scraper:
|
async with CianScraper() as scraper:
|
||||||
scraper.request_delay_sec = request_delay_sec
|
scraper.request_delay_sec = request_delay_sec
|
||||||
|
|
||||||
|
# #1949: per-fetch hard-cancel для browser-fetch'ей внутри bucket-gather.
|
||||||
|
# BrowserFetcher.fetch имеет httpx-timeout 120s, но browser-сервис иногда
|
||||||
|
# виснет так что httpx не получает ответ → gather морозит → heartbeat стоп.
|
||||||
|
# Оборачиваем _browser.fetch в asyncio.wait_for(timeout=X) через monkey-patch:
|
||||||
|
# TimeoutError → _fetch_page_html ловит как Exception → None → _one_page → []
|
||||||
|
# → gather завершается нормально → on_bucket вызывается → heartbeat обновляется.
|
||||||
|
# За флагом cian_full_load_per_fetch_timeout_s (0 = отключить).
|
||||||
|
_fetch_timeout = settings.cian_full_load_per_fetch_timeout_s
|
||||||
|
if _fetch_timeout > 0 and scraper._browser is not None:
|
||||||
|
_orig_fetch = scraper._browser.fetch
|
||||||
|
|
||||||
|
async def _timed_fetch(
|
||||||
|
url: str,
|
||||||
|
_f: Any = _orig_fetch,
|
||||||
|
_t: float = _fetch_timeout,
|
||||||
|
_rid: int = run_id,
|
||||||
|
) -> str:
|
||||||
|
try:
|
||||||
|
return await asyncio.wait_for(_f(url), timeout=_t)
|
||||||
|
except TimeoutError:
|
||||||
|
logger.warning(
|
||||||
|
"cian-full-load run_id=%d: browser-fetch timed out "
|
||||||
|
"after %.0fs — url=%s",
|
||||||
|
_rid,
|
||||||
|
_t,
|
||||||
|
url,
|
||||||
|
)
|
||||||
|
raise
|
||||||
|
|
||||||
|
scraper._browser.fetch = _timed_fetch # type: ignore[method-assign]
|
||||||
|
logger.info(
|
||||||
|
"cian-full-load run_id=%d: per-fetch timeout %.0fs enabled",
|
||||||
|
run_id,
|
||||||
|
_fetch_timeout,
|
||||||
|
)
|
||||||
|
|
||||||
await scraper.fetch_all_secondary(
|
await scraper.fetch_all_secondary(
|
||||||
price_cap_per_bucket=price_cap_per_bucket,
|
price_cap_per_bucket=price_cap_per_bucket,
|
||||||
concurrency=concurrency,
|
concurrency=concurrency,
|
||||||
|
|
@ -2652,6 +2704,15 @@ async def run_cian_full_load(
|
||||||
scrape_runs.mark_failed(db, run_id, str(exc), counters.to_dict())
|
scrape_runs.mark_failed(db, run_id, str(exc), counters.to_dict())
|
||||||
raise
|
raise
|
||||||
|
|
||||||
|
finally:
|
||||||
|
# #1949: отменяем фоновый heartbeat при любом завершении (return / raise / cancel).
|
||||||
|
if _hb_task is not None:
|
||||||
|
_hb_task.cancel()
|
||||||
|
try:
|
||||||
|
await _hb_task
|
||||||
|
except asyncio.CancelledError:
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
# ── Yandex exhaustive full load (region-wide, без anchor'ов) ──────────────────
|
# ── Yandex exhaustive full load (region-wide, без anchor'ов) ──────────────────
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -6,15 +6,26 @@ Fix A (#1950): _rotate_proxy_ip теперь делает до proxy_rotate_atte
|
||||||
Fix B (#1950): run_avito_city_sweep при AvitoBlockedError/RateLimited с
|
Fix B (#1950): run_avito_city_sweep при AvitoBlockedError/RateLimited с
|
||||||
lots_inserted+lots_updated>0 вызывает mark_done (не mark_banned)
|
lots_inserted+lots_updated>0 вызывает mark_done (не mark_banned)
|
||||||
когда avito_serp_ok_not_banned=True.
|
когда avito_serp_ok_not_banned=True.
|
||||||
|
|
||||||
|
Fix C (#1949): run_cian_full_load применяет monkey-patch _browser.fetch с
|
||||||
|
asyncio.wait_for(timeout=cian_full_load_per_fetch_timeout_s).
|
||||||
|
|
||||||
|
Fix D (#1949): run_cian_full_load создаёт фоновый heartbeat-task и отменяет
|
||||||
|
его в finally при любом завершении.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import asyncio
|
||||||
|
import os
|
||||||
|
import time
|
||||||
from typing import Any
|
from typing import Any
|
||||||
from unittest.mock import AsyncMock, MagicMock, patch
|
from unittest.mock import AsyncMock, MagicMock, patch
|
||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
|
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
|
||||||
|
|
||||||
from app.services.scrapers.avito_exceptions import AvitoBlockedError
|
from app.services.scrapers.avito_exceptions import AvitoBlockedError
|
||||||
from app.services.scrapers.base import ScrapedLot
|
from app.services.scrapers.base import ScrapedLot
|
||||||
|
|
||||||
|
|
@ -334,3 +345,214 @@ async def test_avito_city_sweep_flag_off_still_marks_banned() -> None:
|
||||||
mock_runs.mark_banned.called
|
mock_runs.mark_banned.called
|
||||||
), "avito_serp_ok_not_banned=False → mark_banned (backward-compat)"
|
), "avito_serp_ok_not_banned=False → mark_banned (backward-compat)"
|
||||||
assert not mock_runs.mark_done.called
|
assert not mock_runs.mark_done.called
|
||||||
|
|
||||||
|
|
||||||
|
# ──────────────────────────────────────────────────────────────────────────────
|
||||||
|
# Fix C: cian per-fetch wait_for (#1949)
|
||||||
|
# ──────────────────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_cian_full_load_hanging_fetch_cancelled_by_timeout() -> None:
|
||||||
|
"""#1949-C: browser-fetch зависший навсегда отменяется per-fetch timeout.
|
||||||
|
|
||||||
|
Monkey-patch заменяет _browser.fetch на _timed_fetch(wait_for timeout=0.05s).
|
||||||
|
Зависший fetch отменяется через ~50ms → TimeoutError → run завершается.
|
||||||
|
"""
|
||||||
|
from app.services.scrape_pipeline import run_cian_full_load
|
||||||
|
|
||||||
|
fetch_was_called = False
|
||||||
|
|
||||||
|
async def hanging_fetch(url: str) -> str:
|
||||||
|
nonlocal fetch_was_called
|
||||||
|
fetch_was_called = True
|
||||||
|
await asyncio.Event().wait() # висит вечно
|
||||||
|
return ""
|
||||||
|
|
||||||
|
mock_browser = MagicMock()
|
||||||
|
mock_browser.fetch = hanging_fetch
|
||||||
|
|
||||||
|
mock_scraper = MagicMock()
|
||||||
|
mock_scraper.__aenter__ = AsyncMock(return_value=mock_scraper)
|
||||||
|
mock_scraper.__aexit__ = AsyncMock(return_value=None)
|
||||||
|
mock_scraper._browser = mock_browser
|
||||||
|
mock_scraper.request_delay_sec = 0.0
|
||||||
|
|
||||||
|
# fetch_all_secondary симулирует _fetch_page_html: вызывает fetch и ловит Exception
|
||||||
|
async def fake_fetch_all_secondary(**kwargs: object) -> None:
|
||||||
|
try:
|
||||||
|
# После monkey-patch _browser.fetch = _timed_fetch (с timeout=0.05s)
|
||||||
|
await mock_scraper._browser.fetch("https://ekb.cian.ru/sale/flat/1/")
|
||||||
|
except Exception:
|
||||||
|
pass # TimeoutError → аналог _fetch_page_html → returns None
|
||||||
|
|
||||||
|
mock_scraper.fetch_all_secondary = fake_fetch_all_secondary
|
||||||
|
|
||||||
|
mock_db = MagicMock()
|
||||||
|
mock_runs = MagicMock()
|
||||||
|
mock_runs.is_cancelled.return_value = False
|
||||||
|
mock_runs.mark_done = MagicMock()
|
||||||
|
mock_runs.mark_failed = MagicMock()
|
||||||
|
mock_runs.update_heartbeat = MagicMock()
|
||||||
|
|
||||||
|
t0 = time.monotonic()
|
||||||
|
with (
|
||||||
|
patch("app.services.scrapers.cian.CianScraper", return_value=mock_scraper),
|
||||||
|
patch("app.services.scrape_pipeline.scrape_runs", mock_runs),
|
||||||
|
patch("app.services.scrape_pipeline.settings") as s,
|
||||||
|
):
|
||||||
|
s.cian_full_load_per_fetch_timeout_s = 0.05 # 50ms — быстрый timeout для теста
|
||||||
|
s.cian_full_load_detail_top_n = 0
|
||||||
|
|
||||||
|
await run_cian_full_load(mock_db, run_id=1)
|
||||||
|
|
||||||
|
elapsed = time.monotonic() - t0
|
||||||
|
|
||||||
|
assert fetch_was_called, "hanging_fetch должен быть вызван через monkey-patched _browser.fetch"
|
||||||
|
assert (
|
||||||
|
elapsed < 3.0
|
||||||
|
), f"run должен завершиться быстро (fetch timeout 50ms), занял {elapsed:.2f}s"
|
||||||
|
assert mock_runs.mark_done.called, "run_cian_full_load должен завершиться mark_done"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_cian_full_load_per_fetch_timeout_disabled_when_zero() -> None:
|
||||||
|
"""#1949-C: cian_full_load_per_fetch_timeout_s=0 → monkey-patch НЕ применяется."""
|
||||||
|
from app.services.scrape_pipeline import run_cian_full_load
|
||||||
|
|
||||||
|
original_fetch = AsyncMock(return_value="<html/>")
|
||||||
|
|
||||||
|
mock_browser = MagicMock()
|
||||||
|
mock_browser.fetch = original_fetch
|
||||||
|
|
||||||
|
mock_scraper = MagicMock()
|
||||||
|
mock_scraper.__aenter__ = AsyncMock(return_value=mock_scraper)
|
||||||
|
mock_scraper.__aexit__ = AsyncMock(return_value=None)
|
||||||
|
mock_scraper._browser = mock_browser
|
||||||
|
mock_scraper.request_delay_sec = 0.0
|
||||||
|
mock_scraper.fetch_all_secondary = AsyncMock(return_value=None)
|
||||||
|
|
||||||
|
mock_db = MagicMock()
|
||||||
|
mock_runs = MagicMock()
|
||||||
|
mock_runs.is_cancelled.return_value = False
|
||||||
|
mock_runs.mark_done = MagicMock()
|
||||||
|
mock_runs.mark_failed = MagicMock()
|
||||||
|
mock_runs.update_heartbeat = MagicMock()
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch("app.services.scrapers.cian.CianScraper", return_value=mock_scraper),
|
||||||
|
patch("app.services.scrape_pipeline.scrape_runs", mock_runs),
|
||||||
|
patch("app.services.scrape_pipeline.settings") as s,
|
||||||
|
):
|
||||||
|
s.cian_full_load_per_fetch_timeout_s = 0.0 # отключено
|
||||||
|
s.cian_full_load_detail_top_n = 0
|
||||||
|
|
||||||
|
await run_cian_full_load(mock_db, run_id=2)
|
||||||
|
|
||||||
|
# Когда timeout=0 (отключён), monkey-patch не применяется → fetch остаётся original
|
||||||
|
assert (
|
||||||
|
mock_browser.fetch is original_fetch
|
||||||
|
), "cian_full_load_per_fetch_timeout_s=0 → fetch НЕ должен быть заменён monkey-patch"
|
||||||
|
|
||||||
|
|
||||||
|
# ──────────────────────────────────────────────────────────────────────────────
|
||||||
|
# Fix D: cian background heartbeat created and cancelled (#1949)
|
||||||
|
# ──────────────────────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_cian_full_load_background_heartbeat_created_and_cancelled() -> None:
|
||||||
|
"""#1949-D: background heartbeat task создаётся и отменяется в finally."""
|
||||||
|
from app.services.scrape_pipeline import run_cian_full_load
|
||||||
|
|
||||||
|
mock_scraper = MagicMock()
|
||||||
|
mock_scraper.__aenter__ = AsyncMock(return_value=mock_scraper)
|
||||||
|
mock_scraper.__aexit__ = AsyncMock(return_value=None)
|
||||||
|
mock_scraper._browser = None # per-fetch timeout не применяется
|
||||||
|
mock_scraper.fetch_all_secondary = AsyncMock(return_value=None)
|
||||||
|
|
||||||
|
mock_db = MagicMock()
|
||||||
|
mock_runs = MagicMock()
|
||||||
|
mock_runs.is_cancelled.return_value = False
|
||||||
|
mock_runs.mark_done = MagicMock()
|
||||||
|
mock_runs.mark_failed = MagicMock()
|
||||||
|
mock_runs.update_heartbeat = MagicMock()
|
||||||
|
|
||||||
|
created_tasks: list[asyncio.Task[None]] = []
|
||||||
|
_real_create_task = asyncio.create_task
|
||||||
|
|
||||||
|
def _tracking_create_task(coro: Any, *, name: str | None = None) -> asyncio.Task[None]:
|
||||||
|
task = _real_create_task(coro, name=name)
|
||||||
|
created_tasks.append(task)
|
||||||
|
return task # type: ignore[return-value]
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch("app.services.scrapers.cian.CianScraper", return_value=mock_scraper),
|
||||||
|
patch("app.services.scrape_pipeline.scrape_runs", mock_runs),
|
||||||
|
patch(
|
||||||
|
"app.services.scrape_pipeline.asyncio.create_task",
|
||||||
|
side_effect=_tracking_create_task,
|
||||||
|
),
|
||||||
|
patch("app.services.scrape_pipeline.settings") as s,
|
||||||
|
):
|
||||||
|
s.cian_full_load_per_fetch_timeout_s = 0.0
|
||||||
|
s.cian_full_load_detail_top_n = 0
|
||||||
|
|
||||||
|
await run_cian_full_load(mock_db, run_id=3)
|
||||||
|
|
||||||
|
assert (
|
||||||
|
len(created_tasks) >= 1
|
||||||
|
), "asyncio.create_task должен быть вызван хотя бы раз (background heartbeat)"
|
||||||
|
for task in created_tasks:
|
||||||
|
assert (
|
||||||
|
task.cancelled() or task.done()
|
||||||
|
), f"Heartbeat task должен быть отменён в finally: {task}"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_cian_full_load_heartbeat_cancelled_on_exception() -> None:
|
||||||
|
"""#1949-D: heartbeat task отменяется в finally даже при исключении в run."""
|
||||||
|
from app.services.scrape_pipeline import run_cian_full_load
|
||||||
|
|
||||||
|
mock_scraper = MagicMock()
|
||||||
|
mock_scraper.__aenter__ = AsyncMock(return_value=mock_scraper)
|
||||||
|
mock_scraper.__aexit__ = AsyncMock(return_value=None)
|
||||||
|
mock_scraper._browser = None
|
||||||
|
mock_scraper.fetch_all_secondary = AsyncMock(side_effect=RuntimeError("fatal-test-error"))
|
||||||
|
|
||||||
|
mock_db = MagicMock()
|
||||||
|
mock_runs = MagicMock()
|
||||||
|
mock_runs.is_cancelled.return_value = False
|
||||||
|
mock_runs.mark_done = MagicMock()
|
||||||
|
mock_runs.mark_failed = MagicMock()
|
||||||
|
mock_runs.update_heartbeat = MagicMock()
|
||||||
|
|
||||||
|
created_tasks: list[asyncio.Task[None]] = []
|
||||||
|
_real_create_task = asyncio.create_task
|
||||||
|
|
||||||
|
def _tracking_create_task(coro: Any, *, name: str | None = None) -> asyncio.Task[None]:
|
||||||
|
task = _real_create_task(coro, name=name)
|
||||||
|
created_tasks.append(task)
|
||||||
|
return task # type: ignore[return-value]
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch("app.services.scrapers.cian.CianScraper", return_value=mock_scraper),
|
||||||
|
patch("app.services.scrape_pipeline.scrape_runs", mock_runs),
|
||||||
|
patch(
|
||||||
|
"app.services.scrape_pipeline.asyncio.create_task",
|
||||||
|
side_effect=_tracking_create_task,
|
||||||
|
),
|
||||||
|
patch("app.services.scrape_pipeline.settings") as s,
|
||||||
|
):
|
||||||
|
s.cian_full_load_per_fetch_timeout_s = 0.0
|
||||||
|
s.cian_full_load_detail_top_n = 0
|
||||||
|
|
||||||
|
with pytest.raises(RuntimeError, match="fatal-test-error"):
|
||||||
|
await run_cian_full_load(mock_db, run_id=4)
|
||||||
|
|
||||||
|
assert len(created_tasks) >= 1, "heartbeat task должен быть создан"
|
||||||
|
for task in created_tasks:
|
||||||
|
assert (
|
||||||
|
task.cancelled() or task.done()
|
||||||
|
), f"Heartbeat task должен быть отменён в finally даже при RuntimeError: {task}"
|
||||||
|
assert mock_runs.mark_failed.called, "mark_failed должен быть вызван при RuntimeError"
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue