fix(cadastre): remove shared phase_state early-exit (race condition) (#183)
cadastre_jobs.phase_state is shared between ALL parallel Celery workers of the same job. The early-exit triggered as soon as the FIRST worker wrote phase=done via progress_cb — all subsequent workers (reading same shared row) bailed without work. Effect: pilot v7 made 5 NSPD requests for 50 quarters (=0.1/quarter). Full ekb_full job #8: 4950 req / 2408 quarters = 2/quarter — 95% workers early-exited without doing snapshot phase. Root cause of 1.6% parcels coverage. Fix A/B/C only partially helped because most workers never reached Phase 1/1.5. Solution: delete the early-exit. Idempotency is guaranteed by ON CONFLICT DO UPDATE in every upsert_* — re-enqueued tasks write same rows without duplicates. Closes part of #168 (root-cause fix for race condition revealed by #182 metrics). Co-authored-by: lekss361 <claudestars@proton.me>
This commit is contained in:
parent
057f8891dc
commit
e695bca1e4
2 changed files with 50 additions and 32 deletions
|
|
@ -111,23 +111,13 @@ async def harvest_quarter(
|
||||||
"""
|
"""
|
||||||
result = HarvestResult(quarter=quarter)
|
result = HarvestResult(quarter=quarter)
|
||||||
|
|
||||||
# Читаем текущий phase_state из job (для resumability)
|
# NB: `cadastre_jobs.phase_state` шарится между ВСЕМИ parallel Celery workers одного job.
|
||||||
job_row = (
|
# Раньше тут был early-exit `if phase_state.phase == 'done': return` — но он триггерился
|
||||||
db.execute(
|
# как только ПЕРВЫЙ quarter записывал phase=done через update_progress callback, и все
|
||||||
text("SELECT phase_state FROM cadastre_jobs WHERE job_id = :id"),
|
# остальные workers (читая ту же shared row) бейлились без работы. Эффект: 95% quarters
|
||||||
{"id": job_id},
|
# early-exit'или, coverage 1.6% parcels (pilot v7: 5 NSPD req для 50 quarters).
|
||||||
)
|
# Idempotency обеспечивается ON CONFLICT DO UPDATE в каждом upsert_* — re-enqueued task
|
||||||
.mappings()
|
# просто перепишет те же rows, не создавая дубликатов.
|
||||||
.first()
|
|
||||||
)
|
|
||||||
phase_state: dict[str, Any] = {}
|
|
||||||
if job_row and job_row["phase_state"]:
|
|
||||||
phase_state = dict(job_row["phase_state"])
|
|
||||||
|
|
||||||
# Ранний выход если quarter уже полностью обработан
|
|
||||||
if phase_state.get("phase") == "done":
|
|
||||||
logger.info("harvest_quarter: quarter=%s уже done (job_id=%s), skip", quarter, job_id)
|
|
||||||
return result
|
|
||||||
|
|
||||||
# ── Phase 1: search_by_quarter snapshot ──────────────────────────────────
|
# ── Phase 1: search_by_quarter snapshot ──────────────────────────────────
|
||||||
update_progress({"phase": "snapshot_started", "quarter": quarter})
|
update_progress({"phase": "snapshot_started", "quarter": quarter})
|
||||||
|
|
|
||||||
|
|
@ -349,36 +349,64 @@ def test_upsert_parcel_uses_cast_not_double_colon() -> None:
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
async def test_harvest_quarter_phase_done_early_return() -> None:
|
async def test_harvest_quarter_does_not_early_exit_on_shared_phase_done() -> None:
|
||||||
"""Если phase_state='done' — harvest_quarter возвращается ранним выходом."""
|
"""harvest_quarter НЕ должен бейлиться если другой worker записал phase=done
|
||||||
|
в shared `cadastre_jobs.phase_state`. Раньше тут был early-exit, который
|
||||||
|
вызывал 95% workers пропускать работу (race condition между parallel Celery
|
||||||
|
tasks одного job).
|
||||||
|
|
||||||
|
Idempotency обеспечивается ON CONFLICT DO UPDATE в upsert_*, поэтому даже
|
||||||
|
если task re-enqueued после shared phase=done, work должна выполниться.
|
||||||
|
"""
|
||||||
from app.services.cadastre.bulk_harvest import harvest_quarter
|
from app.services.cadastre.bulk_harvest import harvest_quarter
|
||||||
|
|
||||||
|
snapshot = QuarterSnapshot(
|
||||||
|
quarter="66:41:0303161",
|
||||||
|
fetched_at="2026-05-15T10:00:00+00:00",
|
||||||
|
features=[],
|
||||||
|
meta_counts={},
|
||||||
|
)
|
||||||
|
|
||||||
db = MagicMock()
|
db = MagicMock()
|
||||||
|
# Симулируем shared phase_state с phase=done от ДРУГОГО quarter
|
||||||
db.execute = MagicMock(
|
db.execute = MagicMock(
|
||||||
return_value=MagicMock(
|
return_value=MagicMock(
|
||||||
mappings=MagicMock(
|
mappings=MagicMock(
|
||||||
return_value=MagicMock(
|
return_value=MagicMock(
|
||||||
first=MagicMock(return_value={"phase_state": {"phase": "done"}})
|
first=MagicMock(
|
||||||
|
return_value={"phase_state": {"phase": "done", "quarter": "66:41:9999999"}}
|
||||||
|
)
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
|
db.commit = MagicMock()
|
||||||
|
|
||||||
client = AsyncMock()
|
client = AsyncMock()
|
||||||
progress_calls: list[Any] = []
|
client.search_by_quarter = AsyncMock(return_value=snapshot)
|
||||||
|
|
||||||
await harvest_quarter(
|
with patch("app.services.cadastre.bulk_harvest.upsert_features") as mock_upsert:
|
||||||
db=db,
|
mock_upsert.return_value = {
|
||||||
client=client,
|
"parcels": 0,
|
||||||
quarter="66:41:0303161",
|
"buildings": 0,
|
||||||
job_id=1,
|
"constructions": 0,
|
||||||
update_progress=lambda s: progress_calls.append(s),
|
"oncs": 0,
|
||||||
)
|
"enks": 0,
|
||||||
|
"zouit": 0,
|
||||||
|
"skipped": 0,
|
||||||
|
}
|
||||||
|
result = await harvest_quarter(
|
||||||
|
db=db,
|
||||||
|
client=client,
|
||||||
|
quarter="66:41:0303161",
|
||||||
|
job_id=1,
|
||||||
|
update_progress=lambda s: None,
|
||||||
|
)
|
||||||
|
|
||||||
# client.search_by_quarter не должен вызываться
|
# search_by_quarter ДОЛЖЕН быть вызван несмотря на shared phase=done
|
||||||
client.search_by_quarter.assert_not_called()
|
client.search_by_quarter.assert_awaited_once_with("66:41:0303161")
|
||||||
# Нет progress calls
|
assert result.snapshot_requests == 1
|
||||||
assert len(progress_calls) == 0
|
assert result.phase_state == {"phase": "done", "quarter": "66:41:0303161"}
|
||||||
|
|
||||||
|
|
||||||
@pytest.mark.asyncio
|
@pytest.mark.asyncio
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue