fix(tradein/geocode): бюджет прогона проверяется внутри батча (#3151)
All checks were successful
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m52s
Deploy Trade-In / changes (push) Successful in 11s
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / build-backend (push) Successful in 1m7s
Deploy Trade-In / deploy (push) Successful in 1m19s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / perimeter-smoke (push) Successful in 10s
All checks were successful
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / test (push) Successful in 3m52s
Deploy Trade-In / changes (push) Successful in 11s
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / build-backend (push) Successful in 1m7s
Deploy Trade-In / deploy (push) Successful in 1m19s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / perimeter-smoke (push) Successful in 10s
This commit is contained in:
commit
c2d890ff4e
2 changed files with 160 additions and 1 deletions
|
|
@ -56,6 +56,7 @@ class GeocodeBackfillResult:
|
||||||
cache_hits: int = 0 # из geocode_cache (instant)
|
cache_hits: int = 0 # из geocode_cache (instant)
|
||||||
cache_misses: int = 0 # реальные geocoder calls
|
cache_misses: int = 0 # реальные geocoder calls
|
||||||
duration_sec: float = field(default=0.0)
|
duration_sec: float = field(default=0.0)
|
||||||
|
budget_exhausted: bool = False # батч оборван дедлайном, а не разобран до конца
|
||||||
|
|
||||||
|
|
||||||
async def geocode_missing_listings(
|
async def geocode_missing_listings(
|
||||||
|
|
@ -63,6 +64,7 @@ async def geocode_missing_listings(
|
||||||
*,
|
*,
|
||||||
batch_size: int = 200,
|
batch_size: int = 200,
|
||||||
dry_run: bool = False,
|
dry_run: bool = False,
|
||||||
|
deadline_monotonic: float | None = None,
|
||||||
) -> GeocodeBackfillResult:
|
) -> GeocodeBackfillResult:
|
||||||
"""Geocode listings с NULL coords (любой source).
|
"""Geocode listings с NULL coords (любой source).
|
||||||
|
|
||||||
|
|
@ -92,6 +94,8 @@ async def geocode_missing_listings(
|
||||||
Args:
|
Args:
|
||||||
batch_size: max addresses to process per call (default 200 ≈ 3.5 min Nominatim)
|
batch_size: max addresses to process per call (default 200 ≈ 3.5 min Nominatim)
|
||||||
dry_run: только показать что бы сделалось, без UPDATE
|
dry_run: только показать что бы сделалось, без UPDATE
|
||||||
|
deadline_monotonic: значение `time.monotonic()`, после которого батч
|
||||||
|
обрывается на границе адреса (#3151). None — без дедлайна.
|
||||||
|
|
||||||
Returns:
|
Returns:
|
||||||
GeocodeBackfillResult с counters.
|
GeocodeBackfillResult с counters.
|
||||||
|
|
@ -155,6 +159,24 @@ async def geocode_missing_listings(
|
||||||
)
|
)
|
||||||
|
|
||||||
for idx, row in enumerate(rows):
|
for idx, row in enumerate(rows):
|
||||||
|
# Дедлайн проверяется НА КАЖДОМ адресе, а не только между батчами (#3151).
|
||||||
|
# Раньше единственная проверка бюджета стояла в run-обёртке после возврата
|
||||||
|
# из этой функции, то есть потолок прогона на деле был «бюджет + один полный
|
||||||
|
# батч». Пока адрес стоил ~1.9 с это терялось в шуме; после общего
|
||||||
|
# ограничителя темпа Nominatim (#2953) средняя цена 4.5 с, а на трудном
|
||||||
|
# хвосте (tier-1 + до 4 typo-вариантов под паузой 1 с, плюс retry×3) — до
|
||||||
|
# 33 с. Прогон 5017 (27.08) при `budget_sec=1800` шёл 3150 с.
|
||||||
|
# Оборванный батч — штатный исход: `geocode_tried_at` проставлен только у
|
||||||
|
# обработанных пар, необработанные попадут в выборку следующего прогона.
|
||||||
|
if deadline_monotonic is not None and time.monotonic() >= deadline_monotonic:
|
||||||
|
result.budget_exhausted = True
|
||||||
|
logger.info(
|
||||||
|
"geocode_missing: дедлайн исчерпан на %d/%d адресе — обрываю батч",
|
||||||
|
idx,
|
||||||
|
len(rows),
|
||||||
|
)
|
||||||
|
break
|
||||||
|
|
||||||
address: str = row["address"]
|
address: str = row["address"]
|
||||||
city: str | None = row.get("city")
|
city: str | None = row.get("city")
|
||||||
listings_count: int = row["listings_count"]
|
listings_count: int = row["listings_count"]
|
||||||
|
|
@ -356,9 +378,12 @@ async def run_geocode_missing_listings(
|
||||||
runs_mod.update_heartbeat(db, run_id, counters)
|
runs_mod.update_heartbeat(db, run_id, counters)
|
||||||
|
|
||||||
start = time.monotonic()
|
start = time.monotonic()
|
||||||
|
deadline = start + budget_sec
|
||||||
|
|
||||||
while True:
|
while True:
|
||||||
res = await geocode_missing_listings(db, batch_size=batch_size, dry_run=False)
|
res = await geocode_missing_listings(
|
||||||
|
db, batch_size=batch_size, dry_run=False, deadline_monotonic=deadline
|
||||||
|
)
|
||||||
|
|
||||||
# Аккумулируем counters
|
# Аккумулируем counters
|
||||||
total.addresses_total += res.addresses_total
|
total.addresses_total += res.addresses_total
|
||||||
|
|
@ -377,6 +402,18 @@ async def run_geocode_missing_listings(
|
||||||
runs_mod.update_heartbeat(db, run_id, counters)
|
runs_mod.update_heartbeat(db, run_id, counters)
|
||||||
|
|
||||||
elapsed = time.monotonic() - start
|
elapsed = time.monotonic() - start
|
||||||
|
if res.budget_exhausted:
|
||||||
|
# Батч оборван дедлайном внутри себя (#3151) — проверять дренаж
|
||||||
|
# по `addresses_total` уже нельзя: он считает ОТОБРАННЫЕ пары, а не
|
||||||
|
# обработанные, и на оборванном батче ничего не говорит об очереди.
|
||||||
|
logger.info(
|
||||||
|
"run_geocode_missing_listings: run_id=%d — бюджет %.0fs исчерпан "
|
||||||
|
"внутри батча (elapsed=%.1fs), завершаем",
|
||||||
|
run_id,
|
||||||
|
budget_sec,
|
||||||
|
elapsed,
|
||||||
|
)
|
||||||
|
break
|
||||||
if res.addresses_total == 0:
|
if res.addresses_total == 0:
|
||||||
logger.info(
|
logger.info(
|
||||||
"run_geocode_missing_listings: run_id=%d — нет pending адресов, завершаем",
|
"run_geocode_missing_listings: run_id=%d — нет pending адресов, завершаем",
|
||||||
|
|
|
||||||
|
|
@ -2,8 +2,10 @@
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import asyncio
|
||||||
import os
|
import os
|
||||||
import sys
|
import sys
|
||||||
|
import time
|
||||||
from unittest.mock import AsyncMock, MagicMock, patch
|
from unittest.mock import AsyncMock, MagicMock, patch
|
||||||
|
|
||||||
# DATABASE_URL required by config before any app import.
|
# DATABASE_URL required by config before any app import.
|
||||||
|
|
@ -901,3 +903,123 @@ async def test_admin_geocode_missing_select_includes_city_column() -> None:
|
||||||
first_call = db.execute.call_args_list[0]
|
first_call = db.execute.call_args_list[0]
|
||||||
sql_text = str(first_call[0][0])
|
sql_text = str(first_call[0][0])
|
||||||
assert "SELECT id, address, city" in sql_text
|
assert "SELECT id, address, city" in sql_text
|
||||||
|
|
||||||
|
|
||||||
|
# ── #3151: бюджет проверяется ВНУТРИ батча, а не только между батчами ────────
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_geocode_missing_stops_immediately_on_expired_deadline() -> None:
|
||||||
|
"""Дедлайн уже прошёл → ни одного geocode-вызова, батч помечен оборванным."""
|
||||||
|
rows = [
|
||||||
|
{"address": f"ул. Тестовая, {i}", "city": "Екатеринбург", "listings_count": 1}
|
||||||
|
for i in range(5)
|
||||||
|
]
|
||||||
|
db = _mock_db_rows(rows)
|
||||||
|
|
||||||
|
with patch("app.tasks.geocode_missing.geocode", new_callable=AsyncMock) as mock_geo:
|
||||||
|
result = await geocode_missing_listings(
|
||||||
|
db, batch_size=200, deadline_monotonic=time.monotonic() - 1.0
|
||||||
|
)
|
||||||
|
|
||||||
|
mock_geo.assert_not_called()
|
||||||
|
assert result.budget_exhausted is True
|
||||||
|
assert result.addresses_processed == 0
|
||||||
|
assert result.addresses_total == 5 # отобрали 5, обработали 0 — это разные числа
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_geocode_missing_breaks_mid_batch_when_deadline_passes() -> None:
|
||||||
|
"""Дедлайн наступает посреди батча → обрабатываем часть, остальное — следующему прогону."""
|
||||||
|
rows = [
|
||||||
|
{"address": f"ул. Тестовая, {i}", "city": "Екатеринбург", "listings_count": 1}
|
||||||
|
for i in range(6)
|
||||||
|
]
|
||||||
|
db = _mock_db_rows(rows)
|
||||||
|
|
||||||
|
async def _slow_geocode(*_args, **_kwargs):
|
||||||
|
await asyncio.sleep(0.15)
|
||||||
|
return _make_geocode_result()
|
||||||
|
|
||||||
|
with patch("app.tasks.geocode_missing.geocode", side_effect=_slow_geocode):
|
||||||
|
result = await geocode_missing_listings(
|
||||||
|
db, batch_size=200, deadline_monotonic=time.monotonic() + 0.2
|
||||||
|
)
|
||||||
|
|
||||||
|
assert result.budget_exhausted is True
|
||||||
|
assert result.addresses_processed < len(rows)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_geocode_missing_without_deadline_processes_whole_batch() -> None:
|
||||||
|
"""deadline_monotonic=None (умолчание) — поведение прежнее, батч разбирается целиком."""
|
||||||
|
rows = [
|
||||||
|
{"address": f"ул. Тестовая, {i}", "city": "Екатеринбург", "listings_count": 1}
|
||||||
|
for i in range(4)
|
||||||
|
]
|
||||||
|
db = _mock_db_rows(rows)
|
||||||
|
|
||||||
|
with patch(
|
||||||
|
"app.tasks.geocode_missing.geocode",
|
||||||
|
new_callable=AsyncMock,
|
||||||
|
return_value=_make_geocode_result(),
|
||||||
|
) as mock_geo:
|
||||||
|
result = await geocode_missing_listings(db, batch_size=200)
|
||||||
|
|
||||||
|
assert mock_geo.call_count == 4
|
||||||
|
assert result.budget_exhausted is False
|
||||||
|
assert result.addresses_processed == 4
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_run_geocode_missing_listings_stops_on_budget_exhausted_flag() -> None:
|
||||||
|
"""Оборванный по бюджету батч останавливает run-обёртку.
|
||||||
|
|
||||||
|
Регрессия #3151: `addresses_total` на оборванном батче равен размеру ВЫБОРКИ,
|
||||||
|
то есть batch_size — ветка дренажа (`addresses_total < batch_size`) не сработает,
|
||||||
|
а `elapsed > budget_sec` в быстром тесте ложна. Без отдельного флага обёртка
|
||||||
|
крутила бы цикл дальше.
|
||||||
|
"""
|
||||||
|
db = MagicMock()
|
||||||
|
cut_short = GeocodeBackfillResult(
|
||||||
|
addresses_total=200, # == batch_size: дренаж не сработает
|
||||||
|
addresses_processed=12,
|
||||||
|
listings_updated=9,
|
||||||
|
budget_exhausted=True,
|
||||||
|
)
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch(
|
||||||
|
"app.tasks.geocode_missing.geocode_missing_listings",
|
||||||
|
new_callable=AsyncMock,
|
||||||
|
# список, а не постоянный возврат: если обёртка не остановится, тест
|
||||||
|
# упадёт на исчерпании side_effect, а не подвиснет
|
||||||
|
side_effect=[cut_short, cut_short, cut_short],
|
||||||
|
) as mock_batch,
|
||||||
|
patch("app.tasks.geocode_missing.runs_mod") as mock_runs,
|
||||||
|
):
|
||||||
|
result = await run_geocode_missing_listings(db, run_id=11, params={"batch_size": 200})
|
||||||
|
|
||||||
|
assert mock_batch.call_count == 1
|
||||||
|
mock_runs.mark_done.assert_called_once()
|
||||||
|
assert result.addresses_processed == 12
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_run_geocode_missing_listings_passes_deadline_to_batch() -> None:
|
||||||
|
"""Обёртка передаёт в батч дедлайн, а не только считает elapsed после возврата."""
|
||||||
|
db = MagicMock()
|
||||||
|
drained = GeocodeBackfillResult(addresses_total=0)
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch(
|
||||||
|
"app.tasks.geocode_missing.geocode_missing_listings",
|
||||||
|
new_callable=AsyncMock,
|
||||||
|
return_value=drained,
|
||||||
|
) as mock_batch,
|
||||||
|
patch("app.tasks.geocode_missing.runs_mod"),
|
||||||
|
):
|
||||||
|
before = time.monotonic()
|
||||||
|
await run_geocode_missing_listings(db, run_id=12, params={"budget_sec": 90})
|
||||||
|
after = time.monotonic()
|
||||||
|
|
||||||
|
deadline = mock_batch.call_args.kwargs["deadline_monotonic"]
|
||||||
|
assert before + 90 <= deadline <= after + 90
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue