fix(scraper): closing FK violations + transactional poisoning in upsert_objects

ForeignKeyViolation на dev_id='4777_0' / '5377_0' / etc — developer'ов нет
в domrf_developers (FK fk_kn_obj_developer). Plus, старый except: db.rollback()
откатывал ВСЮ транзакцию вместо текущего row → один FK-сбой убивал весь
Phase A.

- _ensure_developers(db, raw) — pre-upsert unique (dev_id, dev_name) перед
  upsert_objects, ON CONFLICT (developer_id) DO NOTHING. Fallback name
  'Developer #<companyGroup>' если все имена None.
- SAVEPOINT (db.begin_nested) во всех 5 upsert helper'ах + _log_failure
  заменяет db.rollback() — теперь сбой одного row не сносит транзакцию.
- Smoke verified против прод-БД: 4777_0/5377_0 закреплены, 6208_0 уже
  существовал — ON CONFLICT отработал.
This commit is contained in:
lekss361 2026-04-28 23:56:27 +03:00
parent d1f2380f32
commit b74c7db438

View file

@ -320,6 +320,24 @@ INSERT_RAW_SQL = text(
) )
def _ensure_snapshot(db: Session, snapshot_date: date) -> None:
"""Гарантирует наличие row в domrf_snapshots для данной даты. Все дочерние
таблицы (domrf_raw_endpoints, domrf_kn_objects, etc) имеют FK на эту
таблицу, поэтому row должен существовать ДО первого INSERT, иначе
ForeignKeyViolation."""
db.execute(
text(
"""
INSERT INTO domrf_snapshots (snapshot_date, fetched_at, notes)
VALUES (:snap, NOW(), 'kn-API sweep')
ON CONFLICT (snapshot_date) DO NOTHING
"""
),
{"snap": snapshot_date},
)
db.commit()
def _insert_raw(db: Session, snapshot_date: date, endpoint_label: str, payload: Any) -> None: def _insert_raw(db: Session, snapshot_date: date, endpoint_label: str, payload: Any) -> None:
body = json.dumps(payload, ensure_ascii=False) body = json.dumps(payload, ensure_ascii=False)
db.execute( db.execute(
@ -335,21 +353,69 @@ def _insert_raw(db: Session, snapshot_date: date, endpoint_label: str, payload:
) )
_UPSERT_DEVELOPER_SQL = text(
"""
INSERT INTO domrf_developers (developer_id, developer_name)
VALUES (:id, :name)
ON CONFLICT (developer_id) DO NOTHING
"""
)
def _ensure_developers(db: Session, raw: list[dict[str, Any]]) -> int:
"""Pre-upsert всех developer'ов, на которых ссылаются объекты — иначе
FK fk_kn_obj_developer (dev_id domrf_developers) валит upsert_objects.
Возвращает кол-во новых записей."""
seen: dict[str, str] = {} # dev_id → name
for row in raw:
dev = row.get("developer") if isinstance(row.get("developer"), dict) else {}
cg = _g(dev, "companyGroup") if dev else None
if not cg:
continue
dev_id = f"{cg}_0"
if dev_id in seen:
continue
# developer_name NOT NULL — fallback на companyGroup как строку.
name = (
_g(dev, "groupName")
or _g(dev, "shortName")
or _g(dev, "fullName")
or f"Developer #{cg}"
)
seen[dev_id] = name
inserted = 0
for dev_id, name in seen.items():
try:
with db.begin_nested():
db.execute(_UPSERT_DEVELOPER_SQL, {"id": dev_id, "name": name})
inserted += 1
except Exception as e:
logger.warning("ensure developer %s failed: %s", dev_id, e)
db.commit()
return inserted
def upsert_objects( def upsert_objects(
db: Session, raw: list[dict[str, Any]], snapshot_date: date, region_cd: int db: Session, raw: list[dict[str, Any]], snapshot_date: date, region_cd: int
) -> int: ) -> int:
# Без этого FK fk_kn_obj_developer валит upsert на новых девелоперах.
_ensure_developers(db, raw)
inserted = 0 inserted = 0
for row in raw: for row in raw:
norm = _norm_object(row, region_cd=region_cd) norm = _norm_object(row, region_cd=region_cd)
if not norm["obj_id"]: if not norm["obj_id"]:
continue continue
norm["snapshot_date"] = snapshot_date norm["snapshot_date"] = snapshot_date
# SAVEPOINT (db.begin_nested) — при FK/CHECK violation одного объекта
# rollback'ается ТОЛЬКО он, остальные uplinks этой транзакции живы.
# Старый db.rollback() сбрасывал всю транзакцию целиком (теряя сотни
# уже успешно закоммиченных объектов).
try: try:
with db.begin_nested():
db.execute(UPSERT_OBJECT_SQL, norm) db.execute(UPSERT_OBJECT_SQL, norm)
inserted += 1 inserted += 1
except Exception as e: except Exception as e:
logger.warning("upsert object %s failed: %s", norm.get("obj_id"), e) logger.warning("upsert object %s failed: %s", norm.get("obj_id"), e)
db.rollback()
return inserted return inserted
@ -365,11 +431,11 @@ def upsert_flats(
continue continue
norm["snapshot_date"] = snapshot_date norm["snapshot_date"] = snapshot_date
try: try:
with db.begin_nested():
db.execute(UPSERT_FLAT_SQL, norm) db.execute(UPSERT_FLAT_SQL, norm)
inserted += 1 inserted += 1
except Exception as e: except Exception as e:
logger.warning("upsert flat %s failed: %s", norm.get("id"), e) logger.warning("upsert flat %s failed: %s", norm.get("id"), e)
db.rollback()
if skipped_no_id: if skipped_no_id:
logger.info("skipped %d flats without numeric id", skipped_no_id) logger.info("skipped %d flats without numeric id", skipped_no_id)
return inserted return inserted
@ -555,6 +621,7 @@ def _log_failure(
if run_id is None: if run_id is None:
return return
try: try:
with db.begin_nested():
db.execute( db.execute(
INSERT_FAILURE_SQL, INSERT_FAILURE_SQL,
{ {
@ -569,7 +636,6 @@ def _log_failure(
) )
except Exception as e: except Exception as e:
logger.warning("failed to log failure: %s", e) logger.warning("failed to log failure: %s", e)
db.rollback()
async def fetch_sale_graph( async def fetch_sale_graph(
@ -616,6 +682,7 @@ def upsert_sale_graph(
if rm is None: if rm is None:
continue continue
try: try:
with db.begin_nested():
db.execute( db.execute(
UPSERT_SALE_GRAPH_SQL, UPSERT_SALE_GRAPH_SQL,
{ {
@ -632,7 +699,6 @@ def upsert_sale_graph(
n += 1 n += 1
except Exception as e: except Exception as e:
logger.warning("upsert sale_graph obj=%s type=%s month=%s: %s", obj_id, type_, rm, e) logger.warning("upsert sale_graph obj=%s type=%s month=%s: %s", obj_id, type_, rm, e)
db.rollback()
return n return n
@ -649,6 +715,7 @@ def upsert_sales_agg(db: Session, obj_id: int, agg: dict[str, Any], snapshot_dat
if not isinstance(d, dict): if not isinstance(d, dict):
continue continue
try: try:
with db.begin_nested():
db.execute( db.execute(
UPSERT_SALES_AGG_SQL, UPSERT_SALES_AGG_SQL,
{ {
@ -664,7 +731,6 @@ def upsert_sales_agg(db: Session, obj_id: int, agg: dict[str, Any], snapshot_dat
n += 1 n += 1
except Exception as e: except Exception as e:
logger.warning("upsert sales_agg obj=%s type=%s: %s", obj_id, type_, e) logger.warning("upsert sales_agg obj=%s type=%s: %s", obj_id, type_, e)
db.rollback()
return n return n
@ -674,6 +740,7 @@ def upsert_infrastructure(
n = 0 n = 0
for poi in pois: for poi in pois:
try: try:
with db.begin_nested():
db.execute( db.execute(
UPSERT_INFRA_SQL, UPSERT_INFRA_SQL,
{ {
@ -691,7 +758,6 @@ def upsert_infrastructure(
n += 1 n += 1
except Exception as e: except Exception as e:
logger.warning("upsert infra obj=%s poi=%s: %s", obj_id, poi.get("name"), e) logger.warning("upsert infra obj=%s poi=%s: %s", obj_id, poi.get("name"), e)
db.rollback()
return n return n
@ -714,6 +780,7 @@ def upsert_photos(
local = local_paths.get(file_id) local = local_paths.get(file_id)
thumb = thumb_paths.get(file_id) thumb = thumb_paths.get(file_id)
try: try:
with db.begin_nested():
db.execute( db.execute(
UPSERT_PHOTO_SQL, UPSERT_PHOTO_SQL,
{ {
@ -736,7 +803,6 @@ def upsert_photos(
n += 1 n += 1
except Exception as e: except Exception as e:
logger.warning("upsert photo obj=%s file=%s: %s", obj_id, file_id, e) logger.warning("upsert photo obj=%s file=%s: %s", obj_id, file_id, e)
db.rollback()
return n return n
@ -968,6 +1034,9 @@ async def run_region_sweep(
""" """
snapshot_date = snapshot_date or date.today() snapshot_date = snapshot_date or date.today()
db = SessionLocal() db = SessionLocal()
# Все дочерние таблицы имеют FK на domrf_snapshots(snapshot_date) —
# гарантируем что row существует ДО любых insert'ов в Phase A/B/C/raw.
_ensure_snapshot(db, snapshot_date)
# ── Init: либо новый run, либо resume ────────────────────────────────── # ── Init: либо новый run, либо resume ──────────────────────────────────
all_objects: list[dict[str, Any]] = [] all_objects: list[dict[str, Any]] = []