Compare commits
No commits in common. "bb74a0e5e10de0d3fb07e5768547417d6e13223b" and "95d348f3c31a151141a9984721d7020c114c36a3" have entirely different histories.
bb74a0e5e1
...
95d348f3c3
3 changed files with 2 additions and 248 deletions
|
|
@ -129,30 +129,10 @@ _SELECT_PENDING_HOUSES = """
|
||||||
SELECT 1 FROM market.yandex_jk_enrichment e WHERE e.ext_id = hs.ext_id
|
SELECT 1 FROM market.yandex_jk_enrichment e WHERE e.ext_id = hs.ext_id
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
-- #2924: дом, у которого резолв уже не удался, не берётся повторно раньше чем
|
ORDER BY h.id, hs.ext_id NULLS LAST
|
||||||
-- через :retry_days (образец — geocode_tried_at, задача geocode_missing).
|
|
||||||
AND (
|
|
||||||
CAST(:force AS boolean) = TRUE
|
|
||||||
OR h.yandex_jk_resolve_tried_at IS NULL
|
|
||||||
OR h.yandex_jk_resolve_tried_at
|
|
||||||
< NOW() - make_interval(days => CAST(:retry_days AS integer))
|
|
||||||
)
|
|
||||||
-- #2924: неудавшиеся — в КОНЕЦ очереди, а не в начало: иначе пять первых по h.id
|
|
||||||
-- занимали все пять слотов каждую неделю.
|
|
||||||
ORDER BY h.id, h.yandex_jk_resolve_tried_at NULLS FIRST, hs.ext_id NULLS LAST
|
|
||||||
LIMIT :lim
|
LIMIT :lim
|
||||||
"""
|
"""
|
||||||
|
|
||||||
# #2924: сколько дней не повторять резолв дома, у которого он не удался. Семь — как у
|
|
||||||
# geocode_tried_at: разметка Яндекса могла поменяться, дом мог появиться в SERP позже.
|
|
||||||
RESOLVE_RETRY_DAYS = 7
|
|
||||||
|
|
||||||
_MARK_RESOLVE_TRIED = """
|
|
||||||
UPDATE houses
|
|
||||||
SET yandex_jk_resolve_tried_at = NOW()
|
|
||||||
WHERE id = CAST(:hid AS bigint)
|
|
||||||
"""
|
|
||||||
|
|
||||||
_UPDATE_SLUG = """
|
_UPDATE_SLUG = """
|
||||||
UPDATE houses
|
UPDATE houses
|
||||||
SET yandex_jk_slug = CAST(:slug AS text)
|
SET yandex_jk_slug = CAST(:slug AS text)
|
||||||
|
|
@ -271,14 +251,7 @@ async def enrich_yandex_newbuilding_sweep(
|
||||||
result.fetchable = sizing["fetchable"]
|
result.fetchable = sizing["fetchable"]
|
||||||
result.pending = sizing["pending"]
|
result.pending = sizing["pending"]
|
||||||
|
|
||||||
rows = (
|
rows = db.execute(text(_SELECT_PENDING_HOUSES), {"force": force, "lim": limit}).mappings().all()
|
||||||
db.execute(
|
|
||||||
text(_SELECT_PENDING_HOUSES),
|
|
||||||
{"force": force, "lim": limit, "retry_days": RESOLVE_RETRY_DAYS},
|
|
||||||
)
|
|
||||||
.mappings()
|
|
||||||
.all()
|
|
||||||
)
|
|
||||||
|
|
||||||
logger.info(
|
logger.info(
|
||||||
"yandex-nb-sweep: total=%d fetchable=%d pending=%d; "
|
"yandex-nb-sweep: total=%d fetchable=%d pending=%d; "
|
||||||
|
|
@ -348,7 +321,6 @@ async def enrich_yandex_newbuilding_sweep(
|
||||||
if not resolved:
|
if not resolved:
|
||||||
logger.warning("slug unresolved house_id=%s ext_id=%s — skip", house_id, ext_id)
|
logger.warning("slug unresolved house_id=%s ext_id=%s — skip", house_id, ext_id)
|
||||||
result.failed_resolve += 1
|
result.failed_resolve += 1
|
||||||
_mark_resolve_tried(db, house_id)
|
|
||||||
continue
|
continue
|
||||||
|
|
||||||
# Persist slug под SAVEPOINT
|
# Persist slug под SAVEPOINT
|
||||||
|
|
@ -363,7 +335,6 @@ async def enrich_yandex_newbuilding_sweep(
|
||||||
"persist slug failed house_id=%s ext_id=%s: %s", house_id, ext_id, exc
|
"persist slug failed house_id=%s ext_id=%s: %s", house_id, ext_id, exc
|
||||||
)
|
)
|
||||||
result.failed_resolve += 1
|
result.failed_resolve += 1
|
||||||
_mark_resolve_tried(db, house_id)
|
|
||||||
continue
|
continue
|
||||||
|
|
||||||
jk_slug = resolved
|
jk_slug = resolved
|
||||||
|
|
@ -477,22 +448,6 @@ async def enrich_yandex_newbuilding_sweep(
|
||||||
return result
|
return result
|
||||||
|
|
||||||
|
|
||||||
def _mark_resolve_tried(db: Session, house_id: int) -> None:
|
|
||||||
"""#2924: зафиксировать неудачную попытку резолва — иначе дом снова первый в очереди.
|
|
||||||
|
|
||||||
Отдельный SAVEPOINT: отметка не должна ронять прогон, а её потеря — не фатальна
|
|
||||||
(дом просто вернётся в очередь на следующем прогоне, как и раньше).
|
|
||||||
"""
|
|
||||||
sp = db.begin_nested()
|
|
||||||
try:
|
|
||||||
db.execute(text(_MARK_RESOLVE_TRIED), {"hid": house_id})
|
|
||||||
sp.commit()
|
|
||||||
db.commit()
|
|
||||||
except Exception as exc:
|
|
||||||
sp.rollback()
|
|
||||||
logger.warning("mark resolve tried failed house_id=%s: %s", house_id, exc)
|
|
||||||
|
|
||||||
|
|
||||||
async def _sleep_with_jitter(delay: float, idx: int, total: int, *, force: bool = False) -> None:
|
async def _sleep_with_jitter(delay: float, idx: int, total: int, *, force: bool = False) -> None:
|
||||||
"""Polite anti-bot sleep с ±20% jitter.
|
"""Polite anti-bot sleep с ±20% jitter.
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1,21 +0,0 @@
|
||||||
-- #2924: неудача резолва yandex_jk_slug ничего не помечала — выборка yandex_newbuilding_sweep
|
|
||||||
-- каждый прогон брала те же дома в том же порядке (ORDER BY h.id), и пять первых
|
|
||||||
-- занимали все пять слотов (limit 5) неделю за неделей. Замер 19.08: 397 ожидающих,
|
|
||||||
-- первые пять по h.id — 6706, 6741, 6757, 6774, 6785 — одни и те же в каждом прогоне.
|
|
||||||
--
|
|
||||||
-- Маркер «пробовали резолвить» — по образцу listings.geocode_tried_at (миграция 005):
|
|
||||||
-- выборка повторяет попытку не раньше чем через N дней и ставит неудавшиеся в конец
|
|
||||||
-- очереди (ORDER BY tried_at NULLS FIRST), а не в начало.
|
|
||||||
--
|
|
||||||
-- Идемпотентно. Колонка nullable, без дефолта — существующие строки = «не пробовали».
|
|
||||||
-- lock_timeout (#2752): ALTER TABLE берёт ACCESS EXCLUSIVE на houses — при живом
|
|
||||||
-- писателе (свип) ждал бы его бесконечно и копил очередь за собой; 5с — и миграция
|
|
||||||
-- падает честно, деплой перезапускается, а не вешает базу.
|
|
||||||
-- SET LOCAL действует только внутри транзакции — вне BEGIN он молча ничего не делает
|
|
||||||
-- (гейт check-migration-lock-timeout это проверяет; на этом я и споткнулся).
|
|
||||||
BEGIN;
|
|
||||||
SET LOCAL lock_timeout = '5s';
|
|
||||||
ALTER TABLE houses ADD COLUMN IF NOT EXISTS yandex_jk_resolve_tried_at timestamptz;
|
|
||||||
COMMENT ON COLUMN houses.yandex_jk_resolve_tried_at IS
|
|
||||||
'#2924: когда yandex_newbuilding_sweep последний раз пытался разрезолвить yandex_jk_slug и не смог; NULL = не пробовали или резолв удался';
|
|
||||||
COMMIT;
|
|
||||||
|
|
@ -1,180 +0,0 @@
|
||||||
"""Неудача резолва yandex_jk_slug помечает дом, а очередь не упирается в те же пять (#2924).
|
|
||||||
|
|
||||||
`_SELECT_PENDING_HOUSES` брал ожидающие дома `ORDER BY h.id` и уходил из очереди только
|
|
||||||
попавший в витрину. Неудачная попытка не писала ничего — ни времени, ни счётчика, — и
|
|
||||||
следующий прогон брал те же дома в том же порядке. Замер 19.08: 397 ожидающих, первые
|
|
||||||
пять по h.id (6706, 6741, 6757, 6774, 6785) занимали все пять слотов каждый прогон.
|
|
||||||
|
|
||||||
Маркер — `houses.yandex_jk_resolve_tried_at` (миграция 269), по образцу
|
|
||||||
`listings.geocode_tried_at`: выборка не берёт дом раньше чем через RESOLVE_RETRY_DAYS,
|
|
||||||
а неудавшиеся идут в КОНЕЦ очереди (`NULLS FIRST`).
|
|
||||||
|
|
||||||
Тесты идут через реальный `enrich_yandex_newbuilding_sweep` с двойником сессии: он отдаёт
|
|
||||||
заданные строки на SELECT и запоминает весь выполненный SQL. Резолвер патчится на None —
|
|
||||||
это «неудача», ради которой маркер и нужен. На origin/main маркер не пишется ни разу.
|
|
||||||
"""
|
|
||||||
|
|
||||||
from __future__ import annotations
|
|
||||||
|
|
||||||
import os
|
|
||||||
|
|
||||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
|
|
||||||
|
|
||||||
from contextlib import contextmanager
|
|
||||||
from typing import Any
|
|
||||||
from unittest.mock import AsyncMock, MagicMock, patch
|
|
||||||
|
|
||||||
import pytest
|
|
||||||
|
|
||||||
|
|
||||||
class _Result:
|
|
||||||
def __init__(self, rows: list[dict], scalar: Any = 0) -> None:
|
|
||||||
self._rows = rows
|
|
||||||
self._scalar = scalar
|
|
||||||
|
|
||||||
def mappings(self) -> _Result:
|
|
||||||
return self
|
|
||||||
|
|
||||||
def all(self) -> list[dict]:
|
|
||||||
return self._rows
|
|
||||||
|
|
||||||
def fetchone(self) -> Any:
|
|
||||||
return self._rows[0] if self._rows else None
|
|
||||||
|
|
||||||
def first(self) -> Any:
|
|
||||||
return self._rows[0] if self._rows else None
|
|
||||||
|
|
||||||
def scalar(self) -> Any:
|
|
||||||
return self._scalar
|
|
||||||
|
|
||||||
def scalar_one(self) -> Any:
|
|
||||||
return self._scalar
|
|
||||||
|
|
||||||
def scalar_one_or_none(self) -> Any:
|
|
||||||
return self._scalar
|
|
||||||
|
|
||||||
|
|
||||||
class _Session:
|
|
||||||
"""Отдаёт pending-строки на SELECT и запоминает все (sql, params)."""
|
|
||||||
|
|
||||||
def __init__(self, pending: list[dict]) -> None:
|
|
||||||
self._pending = pending
|
|
||||||
self.sql: list[tuple[str, dict]] = []
|
|
||||||
|
|
||||||
def execute(self, statement: Any = None, params: Any = None, *a: Any, **kw: Any) -> _Result:
|
|
||||||
text_ = str(statement)
|
|
||||||
self.sql.append((text_, dict(params or {})))
|
|
||||||
# Выборка ожидающих — единственный запрос с DISTINCT ON (h.id); COUNT-запросы
|
|
||||||
# тоже читают houses h + yandex_realty_nb, их путать с выборкой нельзя.
|
|
||||||
if "DISTINCT ON (h.id)" in text_:
|
|
||||||
return _Result(self._pending)
|
|
||||||
return _Result([])
|
|
||||||
|
|
||||||
def begin_nested(self) -> Any:
|
|
||||||
sp = MagicMock()
|
|
||||||
|
|
||||||
@contextmanager
|
|
||||||
def _cm():
|
|
||||||
yield sp
|
|
||||||
|
|
||||||
# код зовёт и `sp = db.begin_nested(); sp.commit()` и `with db.begin_nested():`
|
|
||||||
sp.__enter__ = lambda *_: sp
|
|
||||||
sp.__exit__ = lambda *_: False
|
|
||||||
return sp
|
|
||||||
|
|
||||||
def commit(self) -> None:
|
|
||||||
pass
|
|
||||||
|
|
||||||
def rollback(self) -> None:
|
|
||||||
pass
|
|
||||||
|
|
||||||
|
|
||||||
async def _run(pending: list[dict], *, force: bool = False, resolved: str | None = None):
|
|
||||||
from app.tasks import yandex_newbuilding_sweep as mod
|
|
||||||
|
|
||||||
db = _Session(pending)
|
|
||||||
# Резолвер и конфиг импортируются ВНУТРИ функции (late import) — патчим модуль-источник.
|
|
||||||
with (
|
|
||||||
patch(
|
|
||||||
"scraper_kit.providers.yandex.newbuilding.resolve_yandex_jk_slug",
|
|
||||||
AsyncMock(return_value=resolved),
|
|
||||||
),
|
|
||||||
patch.object(mod, "_sleep_with_jitter", AsyncMock(return_value=None)),
|
|
||||||
patch("app.services.scraper_adapters.RealScraperConfig", MagicMock()),
|
|
||||||
):
|
|
||||||
await mod.enrich_yandex_newbuilding_sweep(
|
|
||||||
db, city="ekaterinburg", limit=5, force=force, request_delay_sec=0
|
|
||||||
)
|
|
||||||
return db
|
|
||||||
|
|
||||||
|
|
||||||
def _tried_updates(db: _Session) -> list[int]:
|
|
||||||
return [
|
|
||||||
int(p["hid"])
|
|
||||||
for s, p in db.sql
|
|
||||||
if "UPDATE houses" in s and "yandex_jk_resolve_tried_at" in s and "hid" in p
|
|
||||||
]
|
|
||||||
|
|
||||||
|
|
||||||
def _select_sql(db: _Session) -> str:
|
|
||||||
for s, _ in db.sql:
|
|
||||||
if "DISTINCT ON (h.id)" in s:
|
|
||||||
return s
|
|
||||||
raise AssertionError("выборка ожидающих домов не выполнялась")
|
|
||||||
|
|
||||||
|
|
||||||
def _select_params(db: _Session) -> dict:
|
|
||||||
for s, p in db.sql:
|
|
||||||
if "DISTINCT ON (h.id)" in s:
|
|
||||||
return p
|
|
||||||
raise AssertionError("выборка ожидающих домов не выполнялась")
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
|
||||||
async def test_failed_resolve_marks_the_house() -> None:
|
|
||||||
"""Головной: резолв не удался → дом получает yandex_jk_resolve_tried_at.
|
|
||||||
|
|
||||||
На origin/main такого UPDATE нет вовсе — дом останется первым в очереди навсегда.
|
|
||||||
"""
|
|
||||||
db = await _run(
|
|
||||||
[{"house_id": 6706, "yandex_jk_slug": None, "yandex_jk_id": None, "ext_id": "286394"}]
|
|
||||||
)
|
|
||||||
assert _tried_updates(db) == [
|
|
||||||
6706
|
|
||||||
], f"неудача резолва не помечена; выполненный SQL: {[s[:50] for s, _ in db.sql]}"
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
|
||||||
async def test_successful_resolve_does_not_mark() -> None:
|
|
||||||
"""Контроль: удачный резолв маркер НЕ ставит — иначе удачные дома тоже уходили
|
|
||||||
бы в конец очереди и ждали бы неделю до обогащения."""
|
|
||||||
db = await _run(
|
|
||||||
[{"house_id": 6706, "yandex_jk_slug": None, "yandex_jk_id": None, "ext_id": "286394"}],
|
|
||||||
resolved="uspenskij",
|
|
||||||
)
|
|
||||||
assert _tried_updates(db) == [], "маркер поставлен при удачном резолве"
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
|
||||||
async def test_selection_skips_recently_tried_and_pushes_tried_to_the_end() -> None:
|
|
||||||
"""Выборка: не берёт дом, у которого попытка свежее RESOLVE_RETRY_DAYS, и ставит
|
|
||||||
пробованные в конец (NULLS FIRST). Проверяется по тексту запроса, который реально
|
|
||||||
ушёл в БД, и по параметрам — retry_days доезжает, а не захардкожен в строке."""
|
|
||||||
db = await _run([])
|
|
||||||
sql = _select_sql(db)
|
|
||||||
assert "yandex_jk_resolve_tried_at IS NULL" in sql, "нет фильтра по маркеру"
|
|
||||||
assert ":retry_days" in sql, "срок повтора не параметризован"
|
|
||||||
assert "yandex_jk_resolve_tried_at NULLS FIRST" in sql, "пробованные не уходят в конец очереди"
|
|
||||||
assert _select_params(db).get("retry_days") == 7, _select_params(db)
|
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
|
||||||
async def test_force_bypasses_the_marker() -> None:
|
|
||||||
"""Контроль: force=True игнорирует маркер — ручной полный проход остаётся возможным."""
|
|
||||||
db = await _run([], force=True)
|
|
||||||
sql = _select_sql(db)
|
|
||||||
# Тот же OR-блок, что и у гейта «уже обогащён»: CAST(:force AS boolean) = TRUE OR …
|
|
||||||
assert (
|
|
||||||
sql.count("CAST(:force AS boolean) = TRUE") >= 2
|
|
||||||
), "маркер не обходится через force — второго OR-блока с :force нет"
|
|
||||||
assert _select_params(db).get("force") is True
|
|
||||||
Loading…
Add table
Reference in a new issue