fix(ptica): успехи батча каталога фиксируются по ходу, а не одним commit'ом в конце (#2464) (#2960)
All checks were successful
Deploy / changes (push) Successful in 8s
Deploy / build-frontend (push) Has been skipped
Deploy / deploy-caddy (push) Has been skipped
Deploy / build-backend (push) Successful in 2m24s
Deploy / build-worker (push) Successful in 3m38s
Deploy / deploy (push) Successful in 1m26s
Deploy / deploy-status (push) Successful in 1s
Deploy / perimeter-smoke (push) Successful in 8s
All checks were successful
Deploy / changes (push) Successful in 8s
Deploy / build-frontend (push) Has been skipped
Deploy / deploy-caddy (push) Has been skipped
Deploy / build-backend (push) Successful in 2m24s
Deploy / build-worker (push) Successful in 3m38s
Deploy / deploy (push) Successful in 1m26s
Deploy / deploy-status (push) Successful in 1s
Deploy / perimeter-smoke (push) Successful in 8s
This commit is contained in:
parent
07b5e5e6da
commit
edbaca1e87
2 changed files with 156 additions and 3 deletions
|
|
@ -462,12 +462,27 @@ async def scrape_catalog_objects(
|
|||
ok = await scrape_catalog_object(db, session, obj_id, snapshot_date)
|
||||
if ok:
|
||||
stats["succeeded"] += 1
|
||||
# Фиксируем сразу, а не одним commit'ом в конце (#2464). Раньше весь
|
||||
# батч жил в одной незакоммиченной транзакции, и любой отказ ПОСЛЕ
|
||||
# цикла — исключение в BrowserSession.__aexit__, снятие Celery-таски,
|
||||
# перезапуск контейнера — обнулял все уже успешные UPDATE'ы.
|
||||
# Это не теория: беговой режим здесь force=True («Загрузить все»),
|
||||
# то есть SQL без LIMIT. На 20.08.2026 в очереди 13200 объектов из
|
||||
# 13801 — многочасовой прогон, где отказ в конце стоил бы всего.
|
||||
# SAVEPOINT внутри scrape_catalog_object к этому моменту уже снят,
|
||||
# поэтому commit здесь корректен.
|
||||
try:
|
||||
db.commit()
|
||||
except Exception:
|
||||
db.rollback()
|
||||
raise
|
||||
else:
|
||||
stats["failed"] += 1
|
||||
|
||||
# Commit outer transaction: SAVEPOINT (`begin_nested`) releases внутри loop,
|
||||
# но outer tx остаётся autobegin'd — без commit() все UPDATE'ы откатятся
|
||||
# при db.close() в Celery task.
|
||||
# Финальный commit. Успешные строки зафиксированы по ходу цикла (см. выше), но
|
||||
# этот вызов остаётся: он закрывает транзакцию, которую могли autobegin'ить
|
||||
# неудачные итерации (их SAVEPOINT откатился, а внешняя транзакция открыта),
|
||||
# и сохраняет прежнее поведение для вызывающих, которые на него полагались.
|
||||
try:
|
||||
db.commit()
|
||||
except Exception:
|
||||
|
|
|
|||
|
|
@ -0,0 +1,138 @@
|
|||
"""Успешные объекты каталога фиксируются по ходу батча, а не одним commit'ом в конце (#2464).
|
||||
|
||||
`scrape_catalog_objects` держал весь батч в одной незакоммиченной транзакции: `db.commit()`
|
||||
стоял ПОСЛЕ блока `async with BrowserSession(...)`. Любой отказ после цикла — исключение в
|
||||
`BrowserSession.__aexit__`, снятие Celery-таски, перезапуск контейнера — обнулял все уже
|
||||
успешные UPDATE'ы.
|
||||
|
||||
Масштаб не гипотетический. Беговой режим здесь `force=True` («Загрузить все»), то есть SQL
|
||||
без LIMIT: на 20.08.2026 в очереди 13200 объектов из 13801 (последний успешный скрейп —
|
||||
19.05, до блокировки DOM.РФ, #2443). Многочасовой прогон, где отказ в конце стоил бы всего.
|
||||
|
||||
Тест не ходит в сеть и в БД: `BrowserSession` и `scrape_catalog_object` подменяются, сессия
|
||||
БД — счётчик вызовов. Проверяется поведение — доживают ли успехи до фиксации, — а не
|
||||
наличие новых символов.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||
|
||||
import asyncio
|
||||
from datetime import date
|
||||
from typing import Any
|
||||
from unittest.mock import patch
|
||||
|
||||
import pytest
|
||||
|
||||
|
||||
class _RecordingSession:
|
||||
"""Сессия-счётчик: помнит порядок commit/rollback."""
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.events: list[str] = []
|
||||
|
||||
def commit(self) -> None:
|
||||
self.events.append("commit")
|
||||
|
||||
def rollback(self) -> None:
|
||||
self.events.append("rollback")
|
||||
|
||||
@property
|
||||
def commits(self) -> int:
|
||||
return self.events.count("commit")
|
||||
|
||||
|
||||
class _FakeBrowserSession:
|
||||
"""Подмена BrowserSession: ничего не делает, но умеет упасть на выходе."""
|
||||
|
||||
raise_on_exit: BaseException | None = None
|
||||
|
||||
def __init__(self, *_a: Any, **_kw: Any) -> None:
|
||||
pass
|
||||
|
||||
async def __aenter__(self) -> _FakeBrowserSession:
|
||||
return self
|
||||
|
||||
async def __aexit__(self, *_exc: Any) -> None:
|
||||
if _FakeBrowserSession.raise_on_exit is not None:
|
||||
exc = _FakeBrowserSession.raise_on_exit
|
||||
_FakeBrowserSession.raise_on_exit = None
|
||||
raise exc
|
||||
|
||||
async def warm_up(self) -> None:
|
||||
return None
|
||||
|
||||
|
||||
def _run(db: _RecordingSession, obj_ids: list[int], *, fail_on_exit: bool = False) -> Any:
|
||||
from app.services.scrapers import domrf_catalog_object as mod
|
||||
|
||||
async def _fake_scrape(_db: Any, _s: Any, obj_id: int, _d: date) -> bool:
|
||||
# Чётные — успех, нечётные — неудача: проверяем, что фиксируются именно успехи.
|
||||
return obj_id % 2 == 0
|
||||
|
||||
_FakeBrowserSession.raise_on_exit = (
|
||||
RuntimeError("падение на выходе из сессии") if fail_on_exit else None
|
||||
)
|
||||
|
||||
with (
|
||||
patch.object(mod, "BrowserSession", _FakeBrowserSession),
|
||||
patch.object(mod, "scrape_catalog_object", _fake_scrape),
|
||||
):
|
||||
return asyncio.run(
|
||||
mod.scrape_catalog_objects(
|
||||
db=db, obj_ids=obj_ids, snapshot_date=date(2026, 8, 20), region_code=66
|
||||
)
|
||||
)
|
||||
|
||||
|
||||
def test_successes_survive_a_late_failure() -> None:
|
||||
"""Отказ ПОСЛЕ цикла не должен стоить уже собранных объектов.
|
||||
|
||||
На origin/main commit стоит после `async with`, поэтому исключение в __aexit__
|
||||
случается ДО единственной фиксации — ни один UPDATE не доживает.
|
||||
"""
|
||||
db = _RecordingSession()
|
||||
|
||||
with pytest.raises(RuntimeError, match="падение на выходе"):
|
||||
_run(db, [2, 4, 6, 8], fail_on_exit=True)
|
||||
|
||||
assert db.commits >= 4, (
|
||||
f"зафиксировано {db.commits} раз при 4 успешных объектах — успехи не пережили "
|
||||
"отказ после цикла: весь батч висел в одной транзакции"
|
||||
)
|
||||
|
||||
|
||||
def test_commit_happens_per_success_not_per_object() -> None:
|
||||
"""Фиксируются успехи, а не каждая итерация — неудачные строки коммитить нечего."""
|
||||
db = _RecordingSession()
|
||||
stats = _run(db, [1, 2, 3, 4, 5, 6])
|
||||
|
||||
assert stats["succeeded"] == 3
|
||||
assert stats["failed"] == 3
|
||||
# 3 успеха по ходу + финальная фиксация в конце.
|
||||
assert db.commits == 4, f"ожидали 3 по ходу + 1 финальную, получили {db.commits}"
|
||||
|
||||
|
||||
def test_empty_batch_touches_nothing() -> None:
|
||||
"""Контроль: пустой список не открывает и не фиксирует ничего."""
|
||||
db = _RecordingSession()
|
||||
stats = _run(db, [])
|
||||
|
||||
assert stats["processed"] == 0
|
||||
assert db.events == [], f"пустой батч тронул сессию: {db.events}"
|
||||
|
||||
|
||||
def test_all_failed_still_closes_transaction() -> None:
|
||||
"""Контроль: батч без единого успеха всё равно закрывает транзакцию.
|
||||
|
||||
Неудачная итерация откатывает свой SAVEPOINT, но внешняя транзакция остаётся
|
||||
autobegin'нутой — финальный commit обязан остаться на месте.
|
||||
"""
|
||||
db = _RecordingSession()
|
||||
stats = _run(db, [1, 3, 5])
|
||||
|
||||
assert stats["succeeded"] == 0
|
||||
assert db.commits == 1, f"ожидали одну финальную фиксацию, получили {db.commits}"
|
||||
Loading…
Add table
Reference in a new issue