gendesign/backend/app/workers/tasks/scrape_objective.py
bot-backend cdf493f345
All checks were successful
Deploy / changes (push) Successful in 9s
Deploy Trade-In / changes (push) Successful in 13s
Deploy / build-frontend (push) Has been skipped
Deploy / deploy-caddy (push) Has been skipped
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy / build-backend (push) Successful in 2m23s
Deploy Trade-In / test (push) Successful in 3m56s
Deploy / build-worker (push) Successful in 4m16s
Deploy Trade-In / build-backend (push) Successful in 1m19s
Deploy / deploy (push) Successful in 1m49s
Deploy / deploy-status (push) Successful in 1s
Deploy / perimeter-smoke (push) Successful in 12s
Deploy Trade-In / deploy (push) Successful in 2m25s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / perimeter-smoke (push) Successful in 11s
chore(format): нормализация под ruff 0.15.20 — 161 файл, только формат (#2864) (#3022)
2026-08-21 12:01:52 +00:00

572 lines
24 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Celery task для регулярного sync с api.objctv.ru.
Стратегия:
Раз в неделю (по умолчанию пн 05:00 МСК) тянем 2 канонических отчёта
по группе из settings.objective_default_group:
- Сводные/Корпуса (за последний полный месяц) → objective_corpus_room_month
- Поквартирные/Лоты (без дат — текущая шахматка) → objective_lots + history
NB: На текущем тарифе Объектива «Сводные/Лоты» и «Поквартирные/Корпуса»
отвечают HTTP 500 (проверено data/sql/69b_objective_minimal.py). Когда тариф
расширят — добавить обратно в `jobs`.
Все payload-ы пишутся в objective_raw_reports (jsonb-колонка) и сразу же
парсятся в нормализованный слой через data/sql/70_parse_objective_raw.py.
Если парсинг упал — raw остаётся, можно re-parse через
`python data/sql/70_parse_objective_raw.py --latest-run`.
"""
from __future__ import annotations
import json
import logging
from datetime import date, timedelta
from typing import Any
from sqlalchemy import text
from sqlalchemy.orm import Session
from app.core.config import settings
from app.core.db import SessionLocal
from app.services.scrapers.objective import (
ObjectiveAPIError,
ObjectiveAuthError,
ObjectiveClient,
)
from app.workers.celery_app import celery_app
logger = logging.getLogger(__name__)
def _start_run(db: Session, group_name: str, triggered_by: str) -> int:
row = db.execute(
text(
"""
INSERT INTO objective_scrape_runs (group_name, triggered_by, status)
VALUES (:gn, :tb, 'running')
RETURNING run_id
"""
),
{"gn": group_name, "tb": triggered_by},
).scalar_one()
db.commit()
return int(row)
def _heartbeat(db: Session, run_id: int, **counts: int) -> None:
sets = ["heartbeat_at = NOW()"]
params: dict[str, Any] = {"rid": run_id}
for k, v in counts.items():
sets.append(f"{k} = :{k}")
params[k] = v
db.execute(
text(f"UPDATE objective_scrape_runs SET {', '.join(sets)} WHERE run_id = :rid"),
params,
)
db.commit()
def _finish_run(
db: Session,
run_id: int,
*,
status: str,
error: str | None = None,
**counts: int,
) -> None:
sets = ["finished_at = NOW()", "status = :status", "error = :error"]
params: dict[str, Any] = {"rid": run_id, "status": status, "error": error}
for k, v in counts.items():
sets.append(f"{k} = :{k}")
params[k] = v
db.execute(
text(f"UPDATE objective_scrape_runs SET {', '.join(sets)} WHERE run_id = :rid"),
params,
)
db.commit()
def _save_raw(
db: Session,
run_id: int,
*,
report_section: str,
report_type: str,
report_name: str,
group_name: str,
complex_name: str | None,
start_date: date | None,
end_date: date | None,
use_ddu: bool,
use_dkp: bool,
payload: Any | None = None,
) -> int:
"""Сохраняет мета + payload в objective_raw_reports.
payload=None допустим для stream-parsed отчётов (lots_pf 600+ МБ).
В этом случае payload_size=0 и payload=NULL (после миграции 79).
"""
if payload is not None:
body = json.dumps(payload, ensure_ascii=False)
payload_param = body
size = len(body.encode("utf-8"))
else:
payload_param = None
size = 0
row = db.execute(
text(
"""
INSERT INTO objective_raw_reports
(run_id, page, report_section, report_type, report_name,
group_name, complex_name, start_date, end_date,
use_ddu, use_dkp, api_version, payload, payload_size)
VALUES (:run_id, 'Отчеты', :section, :type, :name,
:group_name, :complex_name, :start_date, :end_date,
:use_ddu, :use_dkp, 'v2', CAST(:payload AS jsonb), :size)
RETURNING raw_id
"""
),
{
"run_id": run_id,
"section": report_section,
"type": report_type,
"name": report_name,
"group_name": group_name,
"complex_name": complex_name,
"start_date": start_date,
"end_date": end_date,
"use_ddu": use_ddu,
"use_dkp": use_dkp,
"payload": payload_param,
"size": size,
},
).scalar_one()
db.commit()
return int(row)
# Ретраев здесь нет, и параметров, обещающих их, быть не должно (#2464). Стояло
# `bind=True, max_retries=2`, но self не использовался, self.retry() не вызывался и
# autoretry_for задан не был — конфигурация не имела эффекта. Где ретраи нужны, они
# задаются явно: autoretry_for (nspd_sync, scrape_cadastre) или self.retry()
# (scrape_kn.resume_kn_run). Где их сознательно нет — пишется max_retries=0 с
# пояснением (nspd_geo, objective_etl.import_anton_objective).
@celery_app.task(
name="tasks.scrape_objective.sync_objective_group",
)
def sync_objective_group(
group_name: str | None = None,
triggered_by: str = "beat",
use_ddu: bool = True,
use_dkp: bool = True,
period_months_back: int = 1,
rate_ms: int = 3000,
retries: int = 2,
) -> dict[str, Any]:
"""Sync канонические отчёты Объектив для одной группы (территории).
Args:
group_name: имя группы (например 'Свердловская область'). None → settings default
triggered_by: 'beat'/'manual'/'api'/'beat-multi' (журнал)
use_ddu/use_dkp: какие сделки тянуть (передаётся в API)
period_months_back: для Сводных/Корпуса — сколько месяцев назад
захватывать в start_date (1=последний полный месяц)
rate_ms: пауза между HTTP-запросами в client (мс)
retries: ретраи на 429/5xx
Идемпотентность: каждый прогон создаёт новые objective_raw_reports
(history-style). Парсинг сразу через late-import парсера 70.
"""
if not settings.objective_api_key:
logger.warning("sync_objective_group skipped: OBJECTIVE_API_KEY не задан")
return {"skipped": True, "reason": "no_api_key"}
group = group_name or settings.objective_default_group
# period_months_back=1 → последний полный месяц [start_date..end_date)
# period_months_back=N → N последних полных месяцев
end_date = date.today().replace(day=1)
start_date = end_date
for _ in range(max(period_months_back, 1)):
start_date = (start_date - timedelta(days=1)).replace(day=1)
db = SessionLocal()
run_id: int | None = None
try:
run_id = _start_run(db, group, triggered_by)
client = ObjectiveClient(rate_ms=rate_ms, retries=retries)
reports_ok = 0
reports_failed = 0
n_requests = 0
rows_corpus_room = 0
rows_lots = 0
rows_history = 0
# Импортируем парсер и автоматически нормализуем после save_raw.
# Сделано как late-import чтобы scrape_objective.py подгружался без
# data/sql при unit-тестах celery_app.
from importlib import util as _importlib_util
from pathlib import Path as _Path
_parser_path = (
_Path(__file__).resolve().parents[3] / "data" / "sql" / "70_parse_objective_raw.py"
)
_spec = _importlib_util.spec_from_file_location("objective_parser_70", _parser_path)
if _spec is None or _spec.loader is None:
raise RuntimeError(f"Не удалось загрузить парсер по пути {_parser_path}")
parser_mod = _importlib_util.module_from_spec(_spec)
_spec.loader.exec_module(parser_mod)
# На нашем тарифе доступны только Сводные/Корпуса и Поквартирные/Лоты.
# «Сводные/Лоты» и «Поквартирные/Корпуса» возвращают 500 (проверено).
jobs: list[tuple[str, str, dict, str, str, str]] = [
(
"corp_sum",
"report_corpuses_summary",
{
"group_name": group,
"start_date": start_date,
"end_date": end_date,
"use_ddu": use_ddu,
"use_dkp": use_dkp,
},
"Объединенные данные",
"Сводные",
"Корпуса",
),
(
"lots_pf",
"report_lots_per_flat",
{"group_name": group, "use_ddu": use_ddu, "use_dkp": use_dkp},
"Объединенные данные",
"Поквартирные",
"Лоты",
),
]
snap = date.today()
try:
for kind, fn_name, params, section, rtype, rname in jobs:
try:
if kind == "lots_pf":
# lots_pf: 600+ МБ JSON → streaming через ijson, не грузим в RAM.
# payload пишем как NULL в objective_raw_reports (миграция 79).
raw_id = _save_raw(
db,
run_id,
report_section=section,
report_type=rtype,
report_name=rname,
group_name=group,
complex_name=None,
start_date=params.get("start_date"),
end_date=params.get("end_date"),
use_ddu=params.get("use_ddu", True),
use_dkp=params.get("use_dkp", False),
payload=None,
)
# NB: n_requests / reports_ok инкрементируем ТОЛЬКО ПОСЛЕ
# успешного завершения stream+parse+commit. Иначе при провале
# внутреннего generic-except (issue #1220) счётчики бы росли
# даже на упавшем главном отчёте, и _finish_run помечал бы
# run status='done' → silent failure для downstream.
# ObjectiveAuthError/ObjectiveAPIError ⊂ RuntimeError ⊂ Exception
# → ловятся именно этим inner-except (внешний (Auth|API)Error-блок
# на строке ~355 для lots_pf недостижим).
try:
with client.stream_report(
report_type="Поквартирные",
report_name="Лоты",
group_name=group,
use_ddu=params.get("use_ddu", True),
use_dkp=params.get("use_dkp") or False,
) as resp:
n_lots, n_hist = parser_mod.parse_lots_pf_stream(
resp.iter_bytes(chunk_size=65536),
raw_id,
snap,
db,
dry_run=False,
)
rows_lots += n_lots
rows_history += n_hist
db.execute(
text(
"UPDATE objective_raw_reports "
" SET rows_extracted = :n WHERE raw_id = :rid"
),
{"n": n_lots, "rid": raw_id},
)
db.commit()
n_requests += 1
reports_ok += 1
except Exception as parse_err:
db.rollback()
n_requests += 1
reports_failed += 1
logger.exception(
"sync_objective_group: lots_pf stream failed raw_id=%s: %s",
raw_id,
parse_err,
)
else:
# corp_sum (7 МБ) — полный load как раньше
method = getattr(client, fn_name)
payload = method(**params)
n_requests += 1
raw_id = _save_raw(
db,
run_id,
report_section=section,
report_type=rtype,
report_name=rname,
group_name=group,
complex_name=None,
start_date=params.get("start_date"),
end_date=params.get("end_date"),
use_ddu=params.get("use_ddu", True),
use_dkp=params.get("use_dkp", False),
payload=payload,
)
# NB: reports_ok инкрементируем ТОЛЬКО ПОСЛЕ успешного
# parse+commit (зеркально lots_pf / issue #1220, #1369).
# Иначе сбой parse_corp_sum оставлял бы reports_failed==0 →
# _finish_run помечал бы run status='done' → silent failure.
try:
if kind == "corp_sum":
n = parser_mod.parse_corp_sum(
payload, group, raw_id, db, dry_run=False
)
rows_corpus_room += n
db.execute(
text(
"UPDATE objective_raw_reports "
" SET rows_extracted = :n WHERE raw_id = :rid"
),
{"n": n, "rid": raw_id},
)
db.commit()
reports_ok += 1
except Exception as parse_err:
db.rollback()
reports_failed += 1
logger.exception(
"sync_objective_group: parser failed for %s/%s/%s raw_id=%s: %s",
section,
rtype,
rname,
raw_id,
parse_err,
)
_heartbeat(
db,
run_id,
reports_ok=reports_ok,
reports_failed=reports_failed,
requests_count=n_requests,
rows_corpus_room=rows_corpus_room,
rows_lots=rows_lots,
rows_history=rows_history,
)
except (ObjectiveAuthError, ObjectiveAPIError) as e:
reports_failed += 1
logger.warning(
"sync_objective_group: %s/%s/%s failed: %s", section, rtype, rname, e
)
finally:
client.close()
_finish_run(
db,
run_id,
status="done" if reports_failed == 0 else "failed",
error=None if reports_failed == 0 else f"{reports_failed} reports failed",
reports_ok=reports_ok,
reports_failed=reports_failed,
requests_count=n_requests,
rows_corpus_room=rows_corpus_room,
rows_lots=rows_lots,
rows_history=rows_history,
)
return {
"run_id": run_id,
"group_name": group,
"reports_ok": reports_ok,
"reports_failed": reports_failed,
"requests": n_requests,
"rows_corpus_room": rows_corpus_room,
"rows_lots": rows_lots,
"rows_history": rows_history,
"start_date": start_date.isoformat(),
"end_date": end_date.isoformat(),
}
except Exception as e:
if run_id:
# #2464: сессия здесь МОЖЕТ быть отравлена. Исходный сбой бывает
# DB-level (напр. INSERT в _save_raw), и тогда транзакция остаётся в
# aborted-состоянии: следующий execute падает, _finish_run не проходит,
# а голый `except Exception: pass` ниже гасил это молча — строка прогона
# навсегда оставалась в status='running'. Замер прода 20.08: шесть таких
# строк висят с 17.05, то есть 95 суток; уборщика зомби для
# objective_scrape_runs нет.
#
# Сессия здесь СВОЯ (SessionLocal() выше, close в finally), поэтому
# плоский rollback законен: он отбрасывает уже провалившуюся транзакцию
# и ничего чужого не теряет.
try:
db.rollback()
except Exception:
logger.exception("sync_objective_group: rollback перед _finish_run не удался")
try:
_finish_run(db, run_id, status="failed", error=f"{type(e).__name__}: {e}")
except Exception:
# Больше не молча: если и это не прошло, строка останется 'running',
# и знать об этом важнее, чем сохранить тишину в логе.
logger.exception(
"sync_objective_group: не удалось пометить run_id=%s как failed —"
" строка останется в status='running'",
run_id,
)
raise
finally:
db.close()
# ── Multi-group wrapper для еженедельного beat ──────────────────────────────
@celery_app.task(
bind=True,
name="tasks.scrape_objective.sync_all_groups",
max_retries=0,
)
def sync_all_groups(
self: Any,
groups: list[str] | None = None,
triggered_by: str = "beat",
# Все остальные параметры — None = взять из БД-конфига objective_sync_config
use_ddu: bool | None = None,
use_dkp: bool | None = None,
period_months_back: int | None = None,
inter_group_delay_s: int | None = None,
rate_ms: int | None = None,
retries: int | None = None,
) -> dict[str, Any]:
"""Перебирает группы из БД-конфига (или explicit override) и тянет каждую.
Все None-параметры подменяются из таблицы objective_sync_config (PATCH-style).
Это позволяет менять настройки sync через админку без редеплоя.
Args:
groups: explicit override (если None — берётся из БД config.groups_csv)
triggered_by: 'beat' | 'manual' | 'api'
use_ddu/use_dkp/period_months_back/inter_group_delay_s/rate_ms/retries:
None → читается из БД config
Returns:
dict {'groups': [...], 'total_lots': N, 'total_corpus_room': N,
'total_failed': N, 'config_used': {...}}
"""
import time as _time
from app.services.objective_sync_config import get_config
if not settings.objective_api_key:
logger.warning("sync_all_groups skipped: OBJECTIVE_API_KEY не задан")
return {"skipped": True, "reason": "no_api_key"}
# Загружаем динамический config из БД (с fallback на settings)
cfg_db = SessionLocal()
try:
cfg = get_config(cfg_db)
finally:
cfg_db.close()
# PATCH-merge: explicit args > БД config
eff_groups = groups if groups is not None else cfg.groups
eff_use_ddu = use_ddu if use_ddu is not None else cfg.use_ddu
eff_use_dkp = use_dkp if use_dkp is not None else cfg.use_dkp
eff_period = period_months_back if period_months_back is not None else cfg.period_months_back
eff_delay = inter_group_delay_s if inter_group_delay_s is not None else cfg.inter_group_delay_s
eff_rate_ms = rate_ms if rate_ms is not None else cfg.rate_ms
eff_retries = retries if retries is not None else cfg.retries
if not eff_groups:
return {"skipped": True, "reason": "no_groups_configured"}
logger.info(
"sync_all_groups starting for %d groups: %s "
"(use_ddu=%s use_dkp=%s period_months=%s inter_group=%ss rate=%sms retries=%s)",
len(eff_groups),
eff_groups,
eff_use_ddu,
eff_use_dkp,
eff_period,
eff_delay,
eff_rate_ms,
eff_retries,
)
results: list[dict[str, Any]] = []
total_lots = 0
total_crm = 0
total_failed = 0
for idx, group in enumerate(eff_groups):
logger.info("[%d/%d] sync_objective_group(group=%r) START", idx + 1, len(eff_groups), group)
try:
res = sync_objective_group(
group_name=group,
triggered_by=f"{triggered_by}-multi",
use_ddu=eff_use_ddu,
use_dkp=eff_use_dkp,
period_months_back=eff_period,
rate_ms=eff_rate_ms,
retries=eff_retries,
)
if res.get("skipped"):
results.append({"group": group, "status": "skipped", "reason": res.get("reason")})
else:
total_lots += int(res.get("rows_lots") or 0)
total_crm += int(res.get("rows_corpus_room") or 0)
if int(res.get("reports_failed") or 0) > 0:
total_failed += 1
results.append({"group": group, "status": "done", **res})
except Exception as e:
logger.exception(
"[%d/%d] sync for group=%r failed: %s", idx + 1, len(eff_groups), group, e
)
total_failed += 1
results.append(
{"group": group, "status": "failed", "error": f"{type(e).__name__}: {e}"}
)
if idx < len(eff_groups) - 1:
logger.info("inter-group sleep %ds", eff_delay)
_time.sleep(eff_delay)
logger.info(
"sync_all_groups DONE: total_lots=%d total_crm=%d failed=%d",
total_lots,
total_crm,
total_failed,
)
return {
"groups": results,
"total_lots": total_lots,
"total_corpus_room": total_crm,
"total_failed": total_failed,
"n_groups": len(eff_groups),
"config_used": {
"groups": eff_groups,
"use_ddu": eff_use_ddu,
"use_dkp": eff_use_dkp,
"period_months_back": eff_period,
"inter_group_delay_s": eff_delay,
"rate_ms": eff_rate_ms,
"retries": eff_retries,
},
}