feat: bulk-дампы открытых данных ФНС по юрлицам + lookup по ИНН #3429
7 changed files with 1081 additions and 0 deletions
58
tradein-mvp/backend/app/services/fns_lookup.py
Normal file
58
tradein-mvp/backend/app/services/fns_lookup.py
Normal file
|
|
@ -0,0 +1,58 @@
|
||||||
|
"""ФНС opendata lookup по ИНН — тонкое чтение `fns_legal_entity_facts`.
|
||||||
|
|
||||||
|
CONTEXT: см. `fns_opendata_loader.py`. Этот модуль — единственная точка чтения
|
||||||
|
загруженных фактов. ПОТРЕБИТЕЛЯ У НЕГО ПОКА НЕТ: ничто в продукте (скоринг
|
||||||
|
застройщика/УК, estimator и т.п.) сюда не ходит — подключение решается отдельной
|
||||||
|
задачей вне этого PR.
|
||||||
|
|
||||||
|
psycopg v3: SQL через `text(...)` использует `CAST(:x AS type)`, НИКОГДА `:x::type`.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from collections import defaultdict
|
||||||
|
from datetime import date
|
||||||
|
from typing import TypedDict
|
||||||
|
|
||||||
|
from sqlalchemy import text
|
||||||
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
|
_LOOKUP_SQL = text(
|
||||||
|
"""
|
||||||
|
SELECT dataset, series, period, value, org_name
|
||||||
|
FROM fns_legal_entity_facts
|
||||||
|
WHERE inn = CAST(:inn AS text)
|
||||||
|
ORDER BY dataset, series, period
|
||||||
|
"""
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
class FnsFact(TypedDict):
|
||||||
|
series: str
|
||||||
|
period: date
|
||||||
|
value: float
|
||||||
|
org_name: str | None
|
||||||
|
|
||||||
|
|
||||||
|
def get_facts_by_inn(db: Session, inn: str) -> dict[str, list[FnsFact]]:
|
||||||
|
"""Факты по ИНН, сгруппированные по dataset (`revexp`/`sshr2019`/`debtam`/`snr`).
|
||||||
|
|
||||||
|
Пустой/пробельный ИНН и отсутствие данных → `{}` (не исключение — вызывающий код
|
||||||
|
не обязан оборачивать lookup в try/except ради нормального «нет данных»).
|
||||||
|
"""
|
||||||
|
normalized = (inn or "").strip()
|
||||||
|
if not normalized:
|
||||||
|
return {}
|
||||||
|
|
||||||
|
rows = db.execute(_LOOKUP_SQL, {"inn": normalized}).mappings().all()
|
||||||
|
out: dict[str, list[FnsFact]] = defaultdict(list)
|
||||||
|
for row in rows:
|
||||||
|
out[row["dataset"]].append(
|
||||||
|
{
|
||||||
|
"series": row["series"],
|
||||||
|
"period": row["period"],
|
||||||
|
"value": row["value"],
|
||||||
|
"org_name": row["org_name"],
|
||||||
|
}
|
||||||
|
)
|
||||||
|
return dict(out)
|
||||||
497
tradein-mvp/backend/app/services/fns_opendata_loader.py
Normal file
497
tradein-mvp/backend/app/services/fns_opendata_loader.py
Normal file
|
|
@ -0,0 +1,497 @@
|
||||||
|
"""ФНС opendata loader: доходы/расходы, ССЧ, недоимка, спецрежимы юрлиц → lookup по ИНН.
|
||||||
|
|
||||||
|
CONTEXT: открытых данных ЕГРН по правообладателям-физлицам не существует (218-ФЗ
|
||||||
|
ст. 62) — это жёсткий блокер. По ЮРЛИЦАМ данные открыты, лицензия ФНС
|
||||||
|
(nalog.gov.ru/opendata) разрешает переработку и перераспространение — легальный обход
|
||||||
|
для того среза, где он в принципе доступен. Наборы (маска slug'а: 7707329152-<slug>):
|
||||||
|
revexp — доходы и расходы;
|
||||||
|
sshr2019 — среднесписочная численность;
|
||||||
|
debtam — недоимка и задолженность;
|
||||||
|
snr — спецрежимы.
|
||||||
|
ЕГРЮЛ-данных здесь НЕТ: ни адреса, ни ОКВЭД, ни учредителей — только ИНН + наименование
|
||||||
|
+ показатели набора.
|
||||||
|
|
||||||
|
ПОТРЕБИТЕЛЯ У ЭТОГО СЕРВИСА ПОКА НЕТ. Этот модуль (+ fns_lookup.py) — только
|
||||||
|
загрузка и lookup по ИНН. Подключение к продукту (например, скоринг застройщика/УК
|
||||||
|
в оценке) сюда сознательно не входит и не реализовано — это решается отдельной
|
||||||
|
задачей вне текущего PR. `estimator.py` не тронут.
|
||||||
|
|
||||||
|
ИСТОЧНИК: https://www.nalog.gov.ru/opendata/7707329152-<slug>/ — HTML-каталог набора.
|
||||||
|
Прямая ссылка на .zip вида
|
||||||
|
https://file.nalog.ru/opendata/7707329152-revexp/data-20260825-structure-20180110.zip
|
||||||
|
резолвится ДИНАМИЧЕСКИ парсом href со страницы каталога (`resolve_dataset_file_url`):
|
||||||
|
дата в имени файла меняется с каждой публикацией набора — хардкодить URL нельзя,
|
||||||
|
протухнет на следующей публикации. Внутри .zip — множество XML-файлов.
|
||||||
|
|
||||||
|
ПАРСИНГ XML — generic-харвестер, а не хардкод конкретных тегов набора. Точная схема
|
||||||
|
атрибутов XSD structure-20180110 варьируется по набору и НЕ была вживую сверена в
|
||||||
|
этом окружении (нет сетевого доступа к file.nalog.ru отсюда). Вместо этого: любой
|
||||||
|
XML-элемент, несущий ИНН-подобный атрибут (ИННЮЛ/ИНН — данные ФНС opendata, как и
|
||||||
|
ГАР/ФИАС, лежат в атрибутах элементов, не в тексте, см. gar_flats_loader.py), даёт
|
||||||
|
identity строки; ЛЮБОЙ его собственный числовой атрибут (кроме ИНН/КПП/ОКПО/ОКТМО/
|
||||||
|
ОКВЭД/ОГРН и атрибутов-имени) становится отдельным фактом (series=имя атрибута,
|
||||||
|
value=число). Это устойчиво к точным названиям показателей набора ценой чуть более
|
||||||
|
широкого набора series, чем «официальный» словарь показателей — ПЕРЕД первым боевым
|
||||||
|
прогоном на реальном .zip стоит свериться с фактическим XML и, если харвестер тянет
|
||||||
|
лишнее (например служебные коды), сузить `_ID_ATTRS_EXCLUDE`.
|
||||||
|
|
||||||
|
ПАМЯТЬ: .zip целиком лежит в памяти как байты (httpx возвращает `.content` — иначе
|
||||||
|
не проверить целостность архива до распаковки), но НИ ОДИН XML-член НЕ
|
||||||
|
распаковывается в память/на диск целиком: каждый открывается потоково через
|
||||||
|
`zf.open(name)` и парсится `lxml.etree.iterparse` с очисткой обработанных элементов
|
||||||
|
(тот же приём, что в gar_flats_loader.py, — там на многогигабайтных standalone XML,
|
||||||
|
здесь на множестве XML внутри одного архива, открываемых по одному).
|
||||||
|
|
||||||
|
TLS: nalog.gov.ru/file.nalog.ru отдают сертификат НУЦ Минцифры → httpx с default
|
||||||
|
trust store не верифицирует → verify=False. Открытые данные, без auth/PII.
|
||||||
|
Прецеденты: sber_index.py, domrf_kapremont_loader.py.
|
||||||
|
|
||||||
|
psycopg v3: SQL через `text(...)` использует `CAST(:x AS type)`, НИКОГДА `:x::type`.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import io
|
||||||
|
import logging
|
||||||
|
import re
|
||||||
|
import zipfile
|
||||||
|
from collections.abc import Iterator
|
||||||
|
from dataclasses import dataclass
|
||||||
|
from datetime import date
|
||||||
|
from urllib.parse import urljoin
|
||||||
|
|
||||||
|
import httpx
|
||||||
|
from lxml import etree
|
||||||
|
from sqlalchemy import text
|
||||||
|
from sqlalchemy.orm import Session
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
CATALOG_URL_TEMPLATE = "https://www.nalog.gov.ru/opendata/7707329152-{slug}/"
|
||||||
|
DATASET_SLUGS: tuple[str, ...] = ("revexp", "sshr2019", "debtam", "snr")
|
||||||
|
|
||||||
|
DOWNLOAD_TIMEOUT_SEC = 300
|
||||||
|
UPSERT_CHUNK_SIZE = 500
|
||||||
|
|
||||||
|
# Атрибуты-идентификаторы/классификаторы — не показатели, исключаем из харвеста.
|
||||||
|
_INN_ATTRS = ("ИННЮЛ", "ИНН")
|
||||||
|
_NAME_ATTRS = ("НаимОрг", "НаимОрганизации", "НаимОрганизацииПолн", "НаимЮЛПолн")
|
||||||
|
_PERIOD_ATTRS = ("ДатаСост",)
|
||||||
|
# Служебные атрибуты, которые парсятся как числа, но фактом не являются:
|
||||||
|
# идентификаторы, коды и даты документа.
|
||||||
|
_ID_ATTRS_EXCLUDE = frozenset(
|
||||||
|
{
|
||||||
|
"ИННЮЛ",
|
||||||
|
"ИНН",
|
||||||
|
"КПП",
|
||||||
|
"ОКПО",
|
||||||
|
"ОКТМО",
|
||||||
|
"ОКВЭД",
|
||||||
|
"ОГРН",
|
||||||
|
"ИдДок",
|
||||||
|
"ИдФайл",
|
||||||
|
"ДатаДок",
|
||||||
|
"ДатаСост",
|
||||||
|
"ВерсФорм",
|
||||||
|
"ВерсПрог",
|
||||||
|
"КолДок",
|
||||||
|
"ТипИнф",
|
||||||
|
*_NAME_ATTRS,
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
|
_RECORD_TAGS = frozenset({"Документ", "Док", "СвЮЛ", "Сведения"})
|
||||||
|
"""Теги, на которых собирается запись. Реальные выгрузки ФНС используют `Документ`;
|
||||||
|
остальные — запас под соседние наборы. Ограничение обязательно: иначе `Файл`
|
||||||
|
переоткрывал бы уже собранные документы и удваивал записи."""
|
||||||
|
|
||||||
|
_ZIP_HREF_RE = re.compile(r'href="([^"]+\.zip)"', re.IGNORECASE)
|
||||||
|
_PUBLISH_DATE_RE = re.compile(r"data-(\d{8})-")
|
||||||
|
|
||||||
|
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
# Резолв ссылки со страницы каталога (URL меняется на каждой публикации)
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
def resolve_dataset_file_url(html: str, *, base_url: str) -> tuple[str, date | None]:
|
||||||
|
"""Ссылка на .zip актуальной версии набора, снятая со страницы каталога.
|
||||||
|
|
||||||
|
Ищет href, оканчивающийся на .zip и несущий токен даты `data-YYYYMMDD-` — эта
|
||||||
|
дата меняется с каждой публикацией (см. docstring модуля), поэтому URL нельзя
|
||||||
|
хардкодить. Если на странице несколько подходящих ссылок (несколько версий
|
||||||
|
структуры/публикаций), берёт с МАКСИМАЛЬНОЙ датой — детерминированно, без сети.
|
||||||
|
|
||||||
|
Возвращает (абсолютный_url, published_on). published_on — None, если ссылка
|
||||||
|
найдена, но токен даты не распарсился (защитный случай, не должен происходить
|
||||||
|
для ссылок, прошедших регэксп даты).
|
||||||
|
|
||||||
|
Raises ValueError, если на странице нет ни одной ссылки вида *data-YYYYMMDD-*.zip.
|
||||||
|
"""
|
||||||
|
candidates: list[tuple[str, date]] = []
|
||||||
|
for m in _ZIP_HREF_RE.finditer(html):
|
||||||
|
href = m.group(1)
|
||||||
|
date_m = _PUBLISH_DATE_RE.search(href)
|
||||||
|
if not date_m:
|
||||||
|
continue
|
||||||
|
token = date_m.group(1)
|
||||||
|
try:
|
||||||
|
published = date(int(token[:4]), int(token[4:6]), int(token[6:8]))
|
||||||
|
except ValueError:
|
||||||
|
continue
|
||||||
|
candidates.append((href, published))
|
||||||
|
if not candidates:
|
||||||
|
raise ValueError(f"каталог {base_url}: не найдено ссылок вида href=*data-YYYYMMDD-*.zip")
|
||||||
|
href, published = max(candidates, key=lambda c: c[1])
|
||||||
|
return urljoin(base_url, href), published
|
||||||
|
|
||||||
|
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
# In-memory модель факта
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
@dataclass(slots=True)
|
||||||
|
class FnsRecord:
|
||||||
|
inn: str
|
||||||
|
dataset: str
|
||||||
|
series: str
|
||||||
|
period: date
|
||||||
|
value: float
|
||||||
|
org_name: str | None
|
||||||
|
|
||||||
|
def to_params(self) -> dict[str, object]:
|
||||||
|
return {
|
||||||
|
"inn": self.inn,
|
||||||
|
"dataset": self.dataset,
|
||||||
|
"series": self.series,
|
||||||
|
"period": self.period,
|
||||||
|
"value": self.value,
|
||||||
|
"org_name": self.org_name,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
# Стриминговый парсер XML внутри ZIP (lxml iterparse + очистка, без extractall)
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
def _find_attr(elem: etree._Element, candidates: tuple[str, ...]) -> str | None:
|
||||||
|
for attr in candidates:
|
||||||
|
raw = elem.get(attr)
|
||||||
|
if raw:
|
||||||
|
stripped = raw.strip()
|
||||||
|
if stripped:
|
||||||
|
return stripped
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
def _parse_numeric(raw: str) -> float | None:
|
||||||
|
stripped = raw.strip()
|
||||||
|
if not stripped:
|
||||||
|
return None
|
||||||
|
try:
|
||||||
|
return float(stripped.replace(",", "."))
|
||||||
|
except ValueError:
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
def _find_attr_deep(elem: etree._Element, candidates: tuple[str, ...]) -> str | None:
|
||||||
|
"""Найти атрибут на самом элементе ИЛИ на любом его потомке.
|
||||||
|
|
||||||
|
В реальных выгрузках ФНС ИНН и показатели лежат на РАЗНЫХ соседних элементах
|
||||||
|
внутри `<Документ>` — см. докстринг `harvest_element`.
|
||||||
|
"""
|
||||||
|
found = _find_attr(elem, candidates)
|
||||||
|
if found is not None:
|
||||||
|
return found
|
||||||
|
for child in elem.iterdescendants():
|
||||||
|
found = _find_attr(child, candidates)
|
||||||
|
if found is not None:
|
||||||
|
return found
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
def _period_from_element(elem: etree._Element, fallback: date) -> date:
|
||||||
|
"""Отчётная дата документа из `ДатаСост` (ДД.ММ.ГГГГ), иначе — переданная."""
|
||||||
|
raw = _find_attr_deep(elem, _PERIOD_ATTRS)
|
||||||
|
if raw is None:
|
||||||
|
return fallback
|
||||||
|
try:
|
||||||
|
day, month, year = (int(part) for part in raw.split("."))
|
||||||
|
return date(year, month, day)
|
||||||
|
except (ValueError, TypeError):
|
||||||
|
return fallback
|
||||||
|
|
||||||
|
|
||||||
|
def _numeric_attrs(node: etree._Element, seen: set[str]) -> list[tuple[str, float]]:
|
||||||
|
"""Числовые не-служебные атрибуты узла, без повторов по имени серии."""
|
||||||
|
out: list[tuple[str, float]] = []
|
||||||
|
for key, raw in node.attrib.items():
|
||||||
|
if key in _ID_ATTRS_EXCLUDE or key in seen:
|
||||||
|
continue
|
||||||
|
value = _parse_numeric(raw)
|
||||||
|
if value is None:
|
||||||
|
continue
|
||||||
|
seen.add(key)
|
||||||
|
out.append((key, value))
|
||||||
|
return out
|
||||||
|
|
||||||
|
|
||||||
|
def harvest_element(elem: etree._Element, *, dataset: str, period: date) -> list[FnsRecord]:
|
||||||
|
"""Один XML-элемент → 0..N FnsRecord (один на каждый числовой не-id атрибут).
|
||||||
|
|
||||||
|
Реальная форма выгрузки (проверено на data-20260825 набора revexp, 2118 XML
|
||||||
|
в архиве, 40 002 факта на первых 20 001 организаций)::
|
||||||
|
|
||||||
|
<Документ ИдДок="..." ДатаДок="25.08.2026" ДатаСост="31.12.2025">
|
||||||
|
<СведНП НаимОрг="ООО ..." ИННЮЛ="4205406898"/>
|
||||||
|
<СведДохРасх СумДоход="341864000.00" СумРасход="282224000.00"/>
|
||||||
|
</Документ>
|
||||||
|
|
||||||
|
То есть ИНН живёт на `СведНП`, а показатели — на СОСЕДНЕМ элементе. Поэтому
|
||||||
|
носители ИНН ищутся по всему поддереву, а значения без собственного ИНН
|
||||||
|
(`СведДохРасх` и подобные) привязываются к организации ТОЛЬКО когда носитель в
|
||||||
|
поддереве один — иначе непонятно, чьи это цифры, и мы их не выдумываем.
|
||||||
|
|
||||||
|
Отчётный период берётся из `ДатаСост` документа, если он есть; переданный
|
||||||
|
`period` — запасное значение.
|
||||||
|
|
||||||
|
Чистая функция (без БД/сети). Поддерево без ИНН даёт пустой список.
|
||||||
|
"""
|
||||||
|
nodes = [elem, *elem.iterdescendants()]
|
||||||
|
holders = [n for n in nodes if _find_attr(n, _INN_ATTRS) is not None]
|
||||||
|
if not holders:
|
||||||
|
return []
|
||||||
|
doc_period = _period_from_element(elem, period)
|
||||||
|
shared = [n for n in nodes if n not in holders] if len(holders) == 1 else []
|
||||||
|
|
||||||
|
out: list[FnsRecord] = []
|
||||||
|
for holder in holders:
|
||||||
|
inn = _find_attr(holder, _INN_ATTRS)
|
||||||
|
if inn is None: # pragma: no cover - отфильтровано выше
|
||||||
|
continue
|
||||||
|
org_name = _find_attr(holder, _NAME_ATTRS) or (
|
||||||
|
_find_attr_deep(elem, _NAME_ATTRS) if len(holders) == 1 else None
|
||||||
|
)
|
||||||
|
seen: set[str] = set()
|
||||||
|
for node in (holder, *shared):
|
||||||
|
for series, value in _numeric_attrs(node, seen):
|
||||||
|
out.append(
|
||||||
|
FnsRecord(
|
||||||
|
inn=inn,
|
||||||
|
dataset=dataset,
|
||||||
|
series=series,
|
||||||
|
period=doc_period,
|
||||||
|
value=value,
|
||||||
|
org_name=org_name,
|
||||||
|
)
|
||||||
|
)
|
||||||
|
return out
|
||||||
|
|
||||||
|
|
||||||
|
def _iter_xml_records(fh: object, *, dataset: str, period: date) -> Iterator[FnsRecord]:
|
||||||
|
"""Стримит FnsRecord из одного XML-потока `fh` (открытый член ZIP или файл).
|
||||||
|
|
||||||
|
`events=("end",)` + `elem.clear()` + срез предыдущих сиблингов — bounded memory
|
||||||
|
на произвольно большом XML (приём из gar_flats_loader._stream_rows).
|
||||||
|
"""
|
||||||
|
context = etree.iterparse(
|
||||||
|
fh, events=("end",), recover=True, huge_tree=True, resolve_entities=False
|
||||||
|
)
|
||||||
|
for _event, elem in context:
|
||||||
|
if not isinstance(elem.tag, str) or elem.tag not in _RECORD_TAGS:
|
||||||
|
continue
|
||||||
|
yield from harvest_element(elem, dataset=dataset, period=period)
|
||||||
|
# Чистим ТОЛЬКО на границе записи. `end`-события детей приходят раньше
|
||||||
|
# родительского, поэтому безусловный `elem.clear()` на каждом элементе
|
||||||
|
# вычищал `<СведНП>`/`<СведДохРасх>` ДО закрытия `<Документ>` — и парсер
|
||||||
|
# молча извлекал ноль записей из реального дампа.
|
||||||
|
elem.clear()
|
||||||
|
parent = elem.getparent()
|
||||||
|
if parent is not None:
|
||||||
|
while elem.getprevious() is not None:
|
||||||
|
del parent[0]
|
||||||
|
del context
|
||||||
|
|
||||||
|
|
||||||
|
def iter_dataset_records(zip_bytes: bytes, *, dataset: str, period: date) -> Iterator[FnsRecord]:
|
||||||
|
"""Обходит все *.xml внутри `zip_bytes`, элемент за элементом, без extractall.
|
||||||
|
|
||||||
|
Каждый XML-член открывается через `zf.open(name)` как поток (НЕ `zf.read()`
|
||||||
|
целиком, НЕ `zf.extractall()`) — см. docstring модуля § ПАМЯТЬ.
|
||||||
|
Raises ValueError, если в архиве нет ни одного .xml.
|
||||||
|
"""
|
||||||
|
with zipfile.ZipFile(io.BytesIO(zip_bytes)) as zf:
|
||||||
|
names = [n for n in zf.namelist() if n.lower().endswith(".xml")]
|
||||||
|
if not names:
|
||||||
|
raise ValueError(f"zip набора {dataset} не содержит .xml записей: {zf.namelist()!r}")
|
||||||
|
for name in names:
|
||||||
|
with zf.open(name) as fh:
|
||||||
|
yield from _iter_xml_records(fh, dataset=dataset, period=period)
|
||||||
|
|
||||||
|
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
# UPSERT фактов + версия набора
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
_UPSERT_FACTS_SQL = text(
|
||||||
|
"""
|
||||||
|
INSERT INTO fns_legal_entity_facts (inn, dataset, series, period, value, org_name, loaded_at)
|
||||||
|
VALUES (
|
||||||
|
CAST(:inn AS text), CAST(:dataset AS text), CAST(:series AS text),
|
||||||
|
CAST(:period AS date), CAST(:value AS numeric), CAST(:org_name AS text), now()
|
||||||
|
)
|
||||||
|
ON CONFLICT (inn, dataset, series, period) DO UPDATE SET
|
||||||
|
value = EXCLUDED.value,
|
||||||
|
org_name = EXCLUDED.org_name,
|
||||||
|
loaded_at = now()
|
||||||
|
WHERE fns_legal_entity_facts.value IS DISTINCT FROM EXCLUDED.value
|
||||||
|
"""
|
||||||
|
)
|
||||||
|
|
||||||
|
_SELECT_VERSION_SQL = text(
|
||||||
|
"SELECT file_url, published_on FROM fns_dataset_versions WHERE dataset = CAST(:dataset AS text)"
|
||||||
|
)
|
||||||
|
|
||||||
|
_UPSERT_VERSION_SQL = text(
|
||||||
|
"""
|
||||||
|
INSERT INTO fns_dataset_versions (dataset, file_url, published_on, loaded_at)
|
||||||
|
VALUES (CAST(:dataset AS text), CAST(:file_url AS text), CAST(:published_on AS date), now())
|
||||||
|
ON CONFLICT (dataset) DO UPDATE SET
|
||||||
|
file_url = EXCLUDED.file_url,
|
||||||
|
published_on = EXCLUDED.published_on,
|
||||||
|
loaded_at = now()
|
||||||
|
"""
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _chunks(items: list[FnsRecord], size: int) -> Iterator[list[FnsRecord]]:
|
||||||
|
for i in range(0, len(items), size):
|
||||||
|
yield items[i : i + size]
|
||||||
|
|
||||||
|
|
||||||
|
def upsert_records(
|
||||||
|
db: Session, records: list[FnsRecord], *, chunk_size: int = UPSERT_CHUNK_SIZE
|
||||||
|
) -> int:
|
||||||
|
"""Батчевый UPSERT в fns_legal_entity_facts. НЕ коммитит (коммитит caller).
|
||||||
|
|
||||||
|
SAVEPOINT на каждый батч — сбойный батч откатывается изолированно, остальные
|
||||||
|
доезжают (тот же приём, что и gar_flats_loader.upsert_gar_houses).
|
||||||
|
"""
|
||||||
|
upserted = 0
|
||||||
|
failed_batches = 0
|
||||||
|
for batch in _chunks(records, chunk_size):
|
||||||
|
params = [r.to_params() for r in batch]
|
||||||
|
try:
|
||||||
|
with db.begin_nested():
|
||||||
|
db.execute(_UPSERT_FACTS_SQL, params)
|
||||||
|
upserted += len(batch)
|
||||||
|
except Exception:
|
||||||
|
failed_batches += 1
|
||||||
|
logger.warning(
|
||||||
|
"fns_opendata upsert: батч из %d строк сбойнул (пропущен)",
|
||||||
|
len(batch),
|
||||||
|
exc_info=True,
|
||||||
|
)
|
||||||
|
if failed_batches:
|
||||||
|
logger.warning("fns_opendata upsert: сбойных батчей=%d", failed_batches)
|
||||||
|
return upserted
|
||||||
|
|
||||||
|
|
||||||
|
def get_loaded_version(db: Session, dataset: str) -> tuple[str, date | None] | None:
|
||||||
|
"""Уже загруженная версия набора (file_url, published_on) или None."""
|
||||||
|
row = db.execute(_SELECT_VERSION_SQL, {"dataset": dataset}).first()
|
||||||
|
if row is None:
|
||||||
|
return None
|
||||||
|
return row[0], row[1]
|
||||||
|
|
||||||
|
|
||||||
|
def record_dataset_version(
|
||||||
|
db: Session, dataset: str, file_url: str, published_on: date | None
|
||||||
|
) -> None:
|
||||||
|
"""UPSERT версии набора. НЕ коммитит (коммитит caller)."""
|
||||||
|
db.execute(
|
||||||
|
_UPSERT_VERSION_SQL,
|
||||||
|
{"dataset": dataset, "file_url": file_url, "published_on": published_on},
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
# Orchestration
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
def _download(url: str, *, client: httpx.Client) -> bytes:
|
||||||
|
resp = client.get(url, timeout=DOWNLOAD_TIMEOUT_SEC)
|
||||||
|
resp.raise_for_status()
|
||||||
|
return resp.content
|
||||||
|
|
||||||
|
|
||||||
|
def load_dataset(
|
||||||
|
db: Session,
|
||||||
|
slug: str,
|
||||||
|
*,
|
||||||
|
client: httpx.Client | None = None,
|
||||||
|
dry_run: bool = False,
|
||||||
|
force: bool = False,
|
||||||
|
chunk_size: int = UPSERT_CHUNK_SIZE,
|
||||||
|
) -> dict[str, int | str]:
|
||||||
|
"""Резолвит актуальную ссылку набора `slug` → скачивает .zip → парсит → UPSERT.
|
||||||
|
|
||||||
|
Пропускает скачивание, если `fns_dataset_versions` уже содержит ТУ ЖЕ ссылку
|
||||||
|
(force=True форсирует перекачку). dry_run — резолв ссылки происходит (проверяем
|
||||||
|
каталог доступен), скачивание/парс/запись — нет. `client` — для тестов
|
||||||
|
(httpx.MockTransport); если не передан, открывается и закрывается свой.
|
||||||
|
НЕ коммитит — коммитит caller.
|
||||||
|
"""
|
||||||
|
if slug not in DATASET_SLUGS:
|
||||||
|
raise ValueError(f"неизвестный slug набора: {slug!r}, ожидались {DATASET_SLUGS}")
|
||||||
|
|
||||||
|
owns_client = client is None
|
||||||
|
if client is None:
|
||||||
|
# TLS: см. docstring модуля § TLS — verify=False, открытые данные без auth/PII.
|
||||||
|
client = httpx.Client(timeout=DOWNLOAD_TIMEOUT_SEC, verify=False)
|
||||||
|
|
||||||
|
try:
|
||||||
|
catalog_url = CATALOG_URL_TEMPLATE.format(slug=slug)
|
||||||
|
resp = client.get(catalog_url, timeout=DOWNLOAD_TIMEOUT_SEC)
|
||||||
|
resp.raise_for_status()
|
||||||
|
file_url, published_on = resolve_dataset_file_url(resp.text, base_url=catalog_url)
|
||||||
|
|
||||||
|
existing = get_loaded_version(db, slug)
|
||||||
|
if not force and existing is not None and existing[0] == file_url:
|
||||||
|
logger.info("fns_opendata %s: версия %s уже загружена, skip", slug, file_url)
|
||||||
|
return {"dataset": slug, "skipped": 1, "records": 0, "upserted": 0}
|
||||||
|
|
||||||
|
if dry_run:
|
||||||
|
logger.info(
|
||||||
|
"fns_opendata %s: dry_run — резолвлен %s (published_on=%s), скачивание пропущено",
|
||||||
|
slug,
|
||||||
|
file_url,
|
||||||
|
published_on,
|
||||||
|
)
|
||||||
|
return {"dataset": slug, "skipped": 0, "records": 0, "upserted": 0}
|
||||||
|
|
||||||
|
zip_bytes = _download(file_url, client=client)
|
||||||
|
period = published_on or date.today()
|
||||||
|
records = list(iter_dataset_records(zip_bytes, dataset=slug, period=period))
|
||||||
|
upserted = upsert_records(db, records, chunk_size=chunk_size)
|
||||||
|
record_dataset_version(db, slug, file_url, published_on)
|
||||||
|
|
||||||
|
result: dict[str, int | str] = {
|
||||||
|
"dataset": slug,
|
||||||
|
"skipped": 0,
|
||||||
|
"records": len(records),
|
||||||
|
"upserted": upserted,
|
||||||
|
}
|
||||||
|
logger.info("fns_opendata load %s DONE (dry_run=%s): %s", slug, dry_run, result)
|
||||||
|
return result
|
||||||
|
finally:
|
||||||
|
if owns_client:
|
||||||
|
client.close()
|
||||||
|
|
||||||
|
|
||||||
|
def load_datasets(
|
||||||
|
db: Session,
|
||||||
|
slugs: tuple[str, ...] = DATASET_SLUGS,
|
||||||
|
*,
|
||||||
|
client: httpx.Client | None = None,
|
||||||
|
dry_run: bool = False,
|
||||||
|
force: bool = False,
|
||||||
|
) -> dict[str, dict[str, int | str]]:
|
||||||
|
"""load_dataset для каждого slug из `slugs`. НЕ коммитит между наборами (caller)."""
|
||||||
|
return {
|
||||||
|
slug: load_dataset(db, slug, client=client, dry_run=dry_run, force=force) for slug in slugs
|
||||||
|
}
|
||||||
|
|
@ -674,6 +674,53 @@ async def _job_domrf_kapremont_load(
|
||||||
ctx.runs.mark_failed(db, run_id, str(exc)[:1000], {})
|
ctx.runs.mark_failed(db, run_id, str(exc)[:1000], {})
|
||||||
|
|
||||||
|
|
||||||
|
# ── fns_opendata_load — bulk-дампы ФНС по юрлицам → fns_legal_entity_facts ────
|
||||||
|
# Открытые данные ФНС по ЮРЛИЦАМ (лицензия nalog.gov.ru/opendata разрешает
|
||||||
|
# переработку/перераспространение) — легальный обход того, что открытых данных
|
||||||
|
# ЕГРН по правообладателям-физлицам не существует (218-ФЗ ст. 62). Потребителя у
|
||||||
|
# fns_legal_entity_facts на момент добавления НЕТ (см. docstring
|
||||||
|
# app/services/fns_opendata_loader.py и fns_lookup.py) — расписание сидируется
|
||||||
|
# enabled=false, включение отдельным осознанным шагом.
|
||||||
|
async def _job_fns_opendata_load(
|
||||||
|
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||||||
|
) -> None:
|
||||||
|
"""Скачать наборы ФНС opendata (revexp/sshr2019/debtam/snr) → fns_legal_entity_facts.
|
||||||
|
|
||||||
|
Тело переиспользует load_dataset (тот же дизайн-инвариант модуля, что у
|
||||||
|
_job_domrf_kapremont_load: не дублируем логику CLI). commit() после каждого
|
||||||
|
набора — сбой на debtam не должен откатывать уже загруженный revexp.
|
||||||
|
"""
|
||||||
|
from app.services.fns_opendata_loader import DATASET_SLUGS, load_dataset
|
||||||
|
|
||||||
|
datasets = params.get("datasets") or list(DATASET_SLUGS)
|
||||||
|
|
||||||
|
def _run() -> dict[str, int]:
|
||||||
|
total_records = 0
|
||||||
|
total_upserted = 0
|
||||||
|
for slug in datasets:
|
||||||
|
result = load_dataset(db, slug)
|
||||||
|
db.commit()
|
||||||
|
total_records += int(result["records"])
|
||||||
|
total_upserted += int(result["upserted"])
|
||||||
|
return {
|
||||||
|
"records": total_records,
|
||||||
|
"upserted": total_upserted,
|
||||||
|
# см. докстринг _job_domrf_kapremont_load: выделенные колонки прогона +
|
||||||
|
# гейт «три подряд нулевых прогона».
|
||||||
|
"total_seen": total_records,
|
||||||
|
"new_count": total_upserted,
|
||||||
|
}
|
||||||
|
|
||||||
|
loop = asyncio.get_event_loop()
|
||||||
|
try:
|
||||||
|
counters = await loop.run_in_executor(None, _run)
|
||||||
|
ctx.runs.mark_done(db, run_id, counters)
|
||||||
|
except Exception as exc:
|
||||||
|
logger.exception("scheduler: fns_opendata_load crashed run_id=%d", run_id)
|
||||||
|
db.rollback()
|
||||||
|
ctx.runs.mark_failed(db, run_id, str(exc)[:1000], {})
|
||||||
|
|
||||||
|
|
||||||
# ── frt_mkd_load — АИС ППК ФРТ, реестр МКД region 66 (issue #frt-mkd) ────────
|
# ── frt_mkd_load — АИС ППК ФРТ, реестр МКД region 66 (issue #frt-mkd) ────────
|
||||||
async def _job_frt_mkd_load(
|
async def _job_frt_mkd_load(
|
||||||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||||||
|
|
@ -929,6 +976,7 @@ def build_product_handlers(ctx: SchedulerContext) -> dict[str, Handler]:
|
||||||
"domrf_kapremont_load": Handler(_job_domrf_kapremont_load, "domrf_kapremont_load"),
|
"domrf_kapremont_load": Handler(_job_domrf_kapremont_load, "domrf_kapremont_load"),
|
||||||
"frt_mkd_load": Handler(_job_frt_mkd_load, "frt_mkd_load"),
|
"frt_mkd_load": Handler(_job_frt_mkd_load, "frt_mkd_load"),
|
||||||
"cbr_macro_pull": Handler(_job_cbr_macro_pull, "cbr_macro_pull"),
|
"cbr_macro_pull": Handler(_job_cbr_macro_pull, "cbr_macro_pull"),
|
||||||
|
"fns_opendata_load": Handler(_job_fns_opendata_load, "fns_opendata_load"),
|
||||||
"purge_expired_trade_in_data": Handler(
|
"purge_expired_trade_in_data": Handler(
|
||||||
_job_purge_expired_trade_in_data, "purge_expired_trade_in_data"
|
_job_purge_expired_trade_in_data, "purge_expired_trade_in_data"
|
||||||
),
|
),
|
||||||
|
|
|
||||||
71
tradein-mvp/backend/app/tasks/fns_opendata_load.py
Normal file
71
tradein-mvp/backend/app/tasks/fns_opendata_load.py
Normal file
|
|
@ -0,0 +1,71 @@
|
||||||
|
"""CLI: ФНС opendata (revexp/sshr2019/debtam/snr) → fns_legal_entity_facts.
|
||||||
|
|
||||||
|
Только загрузка + lookup (`app/services/fns_lookup.py`) — потребителя у данных пока
|
||||||
|
нет, см. docstring `app/services/fns_opendata_loader.py`.
|
||||||
|
|
||||||
|
Запуск из контейнера tradein-backend/tradein-scraper (нужен интернет к nalog.gov.ru):
|
||||||
|
|
||||||
|
python -m app.tasks.fns_opendata_load # все 4 набора
|
||||||
|
python -m app.tasks.fns_opendata_load --datasets revexp # один набор
|
||||||
|
python -m app.tasks.fns_opendata_load --dry-run # без записи, только резолв ссылки
|
||||||
|
python -m app.tasks.fns_opendata_load --force # перекачать, даже если версия та же
|
||||||
|
|
||||||
|
В режиме --dry-run скачивание/парс/запись не происходят — только резолв актуальной
|
||||||
|
ссылки со страницы каталога (проверка доступности + логирование того, что было бы
|
||||||
|
скачано).
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import argparse
|
||||||
|
import logging
|
||||||
|
|
||||||
|
from app.core.db import SessionLocal
|
||||||
|
from app.services.fns_opendata_loader import DATASET_SLUGS, load_dataset
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
|
def build_parser() -> argparse.ArgumentParser:
|
||||||
|
"""Парсер CLI (вынесен для тестируемости флагов без запуска main)."""
|
||||||
|
parser = argparse.ArgumentParser(
|
||||||
|
description="ФНС opendata loader: revexp/sshr2019/debtam/snr → fns_legal_entity_facts"
|
||||||
|
)
|
||||||
|
parser.add_argument(
|
||||||
|
"--datasets",
|
||||||
|
nargs="+",
|
||||||
|
choices=list(DATASET_SLUGS),
|
||||||
|
default=list(DATASET_SLUGS),
|
||||||
|
help="список наборов (по умолчанию — все 4)",
|
||||||
|
)
|
||||||
|
parser.add_argument(
|
||||||
|
"--dry-run", action="store_true", help="без записи в БД — только резолв ссылки/подсчёт"
|
||||||
|
)
|
||||||
|
parser.add_argument(
|
||||||
|
"--force", action="store_true", help="перекачать, даже если версия совпадает с загруженной"
|
||||||
|
)
|
||||||
|
return parser
|
||||||
|
|
||||||
|
|
||||||
|
def main() -> None:
|
||||||
|
logging.basicConfig(
|
||||||
|
level=logging.INFO,
|
||||||
|
format="%(asctime)s %(levelname)s %(name)s: %(message)s",
|
||||||
|
)
|
||||||
|
args = build_parser().parse_args()
|
||||||
|
|
||||||
|
db = SessionLocal()
|
||||||
|
try:
|
||||||
|
results: dict[str, object] = {}
|
||||||
|
for slug in args.datasets:
|
||||||
|
result = load_dataset(db, slug, dry_run=args.dry_run, force=args.force)
|
||||||
|
if not args.dry_run:
|
||||||
|
db.commit()
|
||||||
|
results[slug] = result
|
||||||
|
logger.info("fns_opendata_load DONE: dry_run=%s %s", args.dry_run, results)
|
||||||
|
finally:
|
||||||
|
db.close()
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
main()
|
||||||
71
tradein-mvp/backend/data/sql/296_fns_legal_entity_facts.sql
Normal file
71
tradein-mvp/backend/data/sql/296_fns_legal_entity_facts.sql
Normal file
|
|
@ -0,0 +1,71 @@
|
||||||
|
-- 296_fns_legal_entity_facts.sql
|
||||||
|
-- ФНС open data по юрлицам: доходы/расходы, среднесписочная численность, недоимка,
|
||||||
|
-- спецрежимы (issue: bulk-дампы ФНС → lookup по ИНН, PR "feat/fns-legal-entities").
|
||||||
|
--
|
||||||
|
-- ПОЧЕМУ ЭТОТ ИСТОЧНИК
|
||||||
|
-- Открытых данных ЕГРН по правообладателям-физлицам не существует (218-ФЗ ст. 62).
|
||||||
|
-- По ЮРЛИЦАМ данные открыты, лицензия ФНС (nalog.gov.ru/opendata) разрешает
|
||||||
|
-- переработку и перераспределение. Наборы (все по маске 7707329152-<slug>):
|
||||||
|
-- revexp (доходы и расходы), sshr2019 (среднесписочная численность), debtam
|
||||||
|
-- (недоимка и задолженность), snr (спецрежимы). ЕГРЮЛ (адрес/ОКВЭД/учредители) в
|
||||||
|
-- этих наборах НЕТ — только ИНН + наименование + показатели набора.
|
||||||
|
--
|
||||||
|
-- ЭТОТ PR — ТОЛЬКО ЗАГРУЗКА И LOOKUP ПО ИНН. Потребителя (скоринг застройщика/УК в
|
||||||
|
-- оценке) здесь нет и не планируется в рамках этого PR — подключение отдельной
|
||||||
|
-- задачей, если решение об этом будет принято отдельно.
|
||||||
|
--
|
||||||
|
-- ФОРМА ХРАНЕНИЯ: одна узкая long-format таблица на все 4 набора, а не таблица на
|
||||||
|
-- набор. Наборы различаются только НАБОРОМ показателей (атрибутов XML), сама форма
|
||||||
|
-- «ИНН + показатель + период + значение» у всех одинаковая — long-format не плодит
|
||||||
|
-- 4 почти идентичные таблицы и не требует миграции при добавлении 5-го набора в
|
||||||
|
-- будущем (не в этом PR).
|
||||||
|
--
|
||||||
|
-- PRIMARY KEY (inn, dataset, series, period), А НЕ (inn, series, period) — так
|
||||||
|
-- предлагал скелет задачи, но series здесь = имя XML-атрибута источника, и без
|
||||||
|
-- dataset в ключе одноимённые атрибуты в двух разных наборах молча схлопнутся в
|
||||||
|
-- одну строку (мы НЕ проверяли по XSD, что имена атрибутов между revexp/sshr2019/
|
||||||
|
-- debtam/snr не пересекаются, поэтому закладываемся на худший случай).
|
||||||
|
--
|
||||||
|
-- Идемпотентно: CREATE TABLE/INDEX IF NOT EXISTS. lock_timeout — на случай, если
|
||||||
|
-- таблицы уже существуют от прогона на другой ветке (ALTER не планируется, но
|
||||||
|
-- дисциплина сохраняется).
|
||||||
|
BEGIN;
|
||||||
|
SET LOCAL lock_timeout = '5s';
|
||||||
|
|
||||||
|
CREATE TABLE IF NOT EXISTS fns_legal_entity_facts (
|
||||||
|
inn text NOT NULL,
|
||||||
|
dataset text NOT NULL, -- 'revexp' | 'sshr2019' | 'debtam' | 'snr'
|
||||||
|
series text NOT NULL, -- имя показателя (атрибут XML источника)
|
||||||
|
period date NOT NULL, -- период показателя; см. fns_opendata_loader.py
|
||||||
|
value numeric NOT NULL,
|
||||||
|
org_name text, -- наименование юрлица на момент загрузки (не ЕГРЮЛ, provenance)
|
||||||
|
loaded_at timestamptz NOT NULL DEFAULT now(),
|
||||||
|
PRIMARY KEY (inn, dataset, series, period)
|
||||||
|
);
|
||||||
|
|
||||||
|
COMMENT ON TABLE fns_legal_entity_facts IS
|
||||||
|
'ФНС opendata (nalog.gov.ru/opendata/7707329152-*) по юрлицам: показатели наборов '
|
||||||
|
'revexp/sshr2019/debtam/snr. Long-format: одна строка = один показатель одного ИНН '
|
||||||
|
'за один период. ЕГРЮЛ-данных (адрес/ОКВЭД/учредители) здесь нет. Без потребителя '
|
||||||
|
'на момент создания — см. app/services/fns_lookup.py.';
|
||||||
|
COMMENT ON COLUMN fns_legal_entity_facts.inn IS 'ИНН юрлица (10 цифр), как в XML источника';
|
||||||
|
COMMENT ON COLUMN fns_legal_entity_facts.dataset IS 'slug набора ФНС opendata: revexp/sshr2019/debtam/snr';
|
||||||
|
COMMENT ON COLUMN fns_legal_entity_facts.series IS 'имя показателя = имя XML-атрибута источника (не переименовываем)';
|
||||||
|
COMMENT ON COLUMN fns_legal_entity_facts.period IS 'период показателя: дата публикации набора (см. fns_dataset_versions.published_on)';
|
||||||
|
COMMENT ON COLUMN fns_legal_entity_facts.value IS 'значение показателя как есть из XML (numeric — источник не различает валюту/единицы в атрибуте)';
|
||||||
|
COMMENT ON COLUMN fns_legal_entity_facts.org_name IS 'наименование юрлица на момент загрузки, provenance only — НЕ авторитетный источник (используй ЕГРЮЛ)';
|
||||||
|
|
||||||
|
CREATE TABLE IF NOT EXISTS fns_dataset_versions (
|
||||||
|
dataset text PRIMARY KEY, -- 'revexp' | 'sshr2019' | 'debtam' | 'snr'
|
||||||
|
file_url text NOT NULL, -- прямая ссылка на .zip, снятая со страницы каталога
|
||||||
|
published_on date, -- дата публикации, распарсенная из имени файла (data-YYYYMMDD-...)
|
||||||
|
loaded_at timestamptz NOT NULL DEFAULT now()
|
||||||
|
);
|
||||||
|
|
||||||
|
COMMENT ON TABLE fns_dataset_versions IS
|
||||||
|
'Последняя загруженная версия каждого набора ФНС opendata — гейт от перекачки '
|
||||||
|
'уже загруженной версии (URL меняется с каждой публикацией, дата в имени файла).';
|
||||||
|
COMMENT ON COLUMN fns_dataset_versions.file_url IS 'резолвится динамически со страницы каталога набора — НЕ хардкодить, протухает';
|
||||||
|
COMMENT ON COLUMN fns_dataset_versions.published_on IS 'дата публикации из имени файла (data-YYYYMMDD-structure-...)';
|
||||||
|
|
||||||
|
COMMIT;
|
||||||
|
|
@ -0,0 +1,36 @@
|
||||||
|
-- 297_scrape_schedules_seed_fns_opendata_load.sql
|
||||||
|
-- Расписание для fns_opendata_load (см. 296_fns_legal_entity_facts.sql — контекст источника).
|
||||||
|
--
|
||||||
|
-- enabled=false: как domclick_detail_backfill (мигр. 175) — включение отдельным
|
||||||
|
-- осознанным шагом после деплоя, дымовой пробы и ручного --dry-run прогона. Наборы
|
||||||
|
-- ФНС весят десятки-сотни МБ каждый, первый прогон надо смотреть глазами.
|
||||||
|
--
|
||||||
|
-- interval_days=30: наборы ФНС публикуются раз в квартал/год (revexp — по итогам
|
||||||
|
-- отчётного периода), суточный/недельный опрос каталога бессмысленно частый.
|
||||||
|
-- Реализовано через default_params (interval_minutes читает post_claim
|
||||||
|
-- reschedule_after_minutes, здесь используем ту же схему в днях: 30 * 24 * 60).
|
||||||
|
--
|
||||||
|
-- ЗАВИСИМОСТИ: 052_scrape_schedules.sql (таблица + UNIQUE(source)), 294 (таблицы факта).
|
||||||
|
-- Идемпотентно: ON CONFLICT (source) DO NOTHING.
|
||||||
|
BEGIN;
|
||||||
|
|
||||||
|
INSERT INTO scrape_schedules (
|
||||||
|
source,
|
||||||
|
enabled,
|
||||||
|
window_start_hour,
|
||||||
|
window_end_hour,
|
||||||
|
next_run_at,
|
||||||
|
default_params
|
||||||
|
)
|
||||||
|
VALUES
|
||||||
|
(
|
||||||
|
'fns_opendata_load',
|
||||||
|
false,
|
||||||
|
2,
|
||||||
|
5,
|
||||||
|
((CURRENT_DATE + INTERVAL '1 day')) AT TIME ZONE 'UTC',
|
||||||
|
'{"interval_minutes": 43200, "datasets": ["revexp", "sshr2019", "debtam", "snr"]}'::jsonb
|
||||||
|
)
|
||||||
|
ON CONFLICT (source) DO NOTHING;
|
||||||
|
|
||||||
|
COMMIT;
|
||||||
300
tradein-mvp/backend/tests/test_fns_opendata_loader.py
Normal file
300
tradein-mvp/backend/tests/test_fns_opendata_loader.py
Normal file
|
|
@ -0,0 +1,300 @@
|
||||||
|
"""Тесты ФНС opendata loader'а + lookup (app/services/fns_opendata_loader.py,
|
||||||
|
app/services/fns_lookup.py, мигр. 296/295, "bulk-дампы ФНС по юрлицам").
|
||||||
|
|
||||||
|
Coverage:
|
||||||
|
- resolve_dataset_file_url — резолв .zip-ссылки из литерального куска HTML каталога
|
||||||
|
(несколько ссылок → берём максимальную дату публикации); ValueError, если ссылок нет.
|
||||||
|
- harvest_element — generic-харвестер: ИНН-подобный атрибут даёт identity, числовые
|
||||||
|
атрибуты (кроме id/name-подобных) становятся FnsRecord; элемент без ИНН → [].
|
||||||
|
- iter_dataset_records — round-trip через реально собранный .zip с мини-XML фикстурой.
|
||||||
|
- load_dataset — HTTP замокан через httpx.MockTransport (каталог + zip); dry_run не
|
||||||
|
пишет; повторная загрузка той же версии — skip без сети на .zip.
|
||||||
|
- Статические asserts по SQL: CAST-дисциплина, ON CONFLICT, IS DISTINCT FROM.
|
||||||
|
- upsert_records: `db.begin_nested()` в исходнике, модуль НЕ коммитит (коммитит caller).
|
||||||
|
- fns_lookup.get_facts_by_inn на MagicMock db: группировка по dataset, пустой ИНН.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import inspect
|
||||||
|
import io
|
||||||
|
import os
|
||||||
|
import re
|
||||||
|
import zipfile
|
||||||
|
from datetime import date
|
||||||
|
from unittest.mock import MagicMock
|
||||||
|
|
||||||
|
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||||||
|
|
||||||
|
import httpx
|
||||||
|
import pytest
|
||||||
|
from lxml import etree
|
||||||
|
|
||||||
|
from app.services import fns_lookup as fl
|
||||||
|
from app.services import fns_opendata_loader as fol
|
||||||
|
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
# resolve_dataset_file_url
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
_CATALOG_HTML = """
|
||||||
|
<html><body>
|
||||||
|
<table>
|
||||||
|
<tr><td><a href="https://file.nalog.ru/opendata/7707329152-revexp/data-20260101-structure-20180110.zip">2026-01-01</a></td></tr>
|
||||||
|
<tr><td><a href="https://file.nalog.ru/opendata/7707329152-revexp/data-20260825-structure-20180110.zip">2026-08-25</a></td></tr>
|
||||||
|
<tr><td><a href="/opendata/7707329152-revexp/structure-20180110.xsd">XSD</a></td></tr>
|
||||||
|
</table>
|
||||||
|
</body></html>
|
||||||
|
"""
|
||||||
|
|
||||||
|
|
||||||
|
def test_resolve_dataset_file_url_picks_latest_publication() -> None:
|
||||||
|
url, published = fol.resolve_dataset_file_url(
|
||||||
|
_CATALOG_HTML, base_url="https://www.nalog.gov.ru/opendata/7707329152-revexp/"
|
||||||
|
)
|
||||||
|
assert (
|
||||||
|
url
|
||||||
|
== "https://file.nalog.ru/opendata/7707329152-revexp/data-20260825-structure-20180110.zip"
|
||||||
|
)
|
||||||
|
assert published == date(2026, 8, 25)
|
||||||
|
|
||||||
|
|
||||||
|
def test_resolve_dataset_file_url_relative_href_resolved_against_base() -> None:
|
||||||
|
html = '<a href="data-20260301-structure-1.zip">v</a>'
|
||||||
|
url, published = fol.resolve_dataset_file_url(
|
||||||
|
html, base_url="https://www.nalog.gov.ru/opendata/7707329152-snr/"
|
||||||
|
)
|
||||||
|
assert url == "https://www.nalog.gov.ru/opendata/7707329152-snr/data-20260301-structure-1.zip"
|
||||||
|
assert published == date(2026, 3, 1)
|
||||||
|
|
||||||
|
|
||||||
|
def test_resolve_dataset_file_url_raises_when_no_zip_links() -> None:
|
||||||
|
with pytest.raises(ValueError, match="data-YYYYMMDD"):
|
||||||
|
fol.resolve_dataset_file_url("<html><body>empty</body></html>", base_url="https://x/")
|
||||||
|
|
||||||
|
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
# harvest_element — generic-харвестер
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
def test_harvest_element_extracts_numeric_attrs_excluding_ids() -> None:
|
||||||
|
elem = etree.fromstring(
|
||||||
|
'<СвНП ИННЮЛ="7707329152" КПП="770701001" НаимОрг="ООО РОМАШКА" '
|
||||||
|
'СумДоход="1234567.89" СумРасход="654321,00"/>'
|
||||||
|
)
|
||||||
|
records = fol.harvest_element(elem, dataset="revexp", period=date(2026, 8, 25))
|
||||||
|
by_series = {r.series: r for r in records}
|
||||||
|
assert set(by_series) == {"СумДоход", "СумРасход"}
|
||||||
|
assert by_series["СумДоход"].inn == "7707329152"
|
||||||
|
assert by_series["СумДоход"].value == pytest.approx(1234567.89)
|
||||||
|
assert by_series["СумРасход"].value == pytest.approx(654321.00) # запятая — decimal separator
|
||||||
|
assert by_series["СумДоход"].org_name == "ООО РОМАШКА"
|
||||||
|
assert by_series["СумДоход"].dataset == "revexp"
|
||||||
|
assert by_series["СумДоход"].period == date(2026, 8, 25)
|
||||||
|
|
||||||
|
|
||||||
|
def test_harvest_element_without_inn_returns_empty() -> None:
|
||||||
|
elem = etree.fromstring('<Прочее КПП="770701001" Значение="1"/>')
|
||||||
|
assert fol.harvest_element(elem, dataset="revexp", period=date(2026, 1, 1)) == []
|
||||||
|
|
||||||
|
|
||||||
|
def test_harvest_element_skips_non_numeric_attrs() -> None:
|
||||||
|
elem = etree.fromstring('<СвНП ИННЮЛ="123" Комментарий="текст, не число"/>')
|
||||||
|
assert fol.harvest_element(elem, dataset="snr", period=date(2026, 1, 1)) == []
|
||||||
|
|
||||||
|
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
# iter_dataset_records — round-trip через собранный .zip
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
def _zip_bytes(xml_by_name: dict[str, str]) -> bytes:
|
||||||
|
buf = io.BytesIO()
|
||||||
|
with zipfile.ZipFile(buf, "w") as zf:
|
||||||
|
for name, content in xml_by_name.items():
|
||||||
|
zf.writestr(name, content)
|
||||||
|
return buf.getvalue()
|
||||||
|
|
||||||
|
|
||||||
|
_SAMPLE_XML = """<?xml version="1.0" encoding="UTF-8"?>
|
||||||
|
<Файл>
|
||||||
|
<Документ>
|
||||||
|
<СвНП ИННЮЛ="7707329152" НаимОрг="ООО РОМАШКА" СумДоход="1000"/>
|
||||||
|
<СвНП ИННЮЛ="6660000001" НаимОрг="ООО ВАСИЛЁК" СумДоход="2000" СумРасход="500"/>
|
||||||
|
</Документ>
|
||||||
|
</Файл>
|
||||||
|
"""
|
||||||
|
|
||||||
|
|
||||||
|
def test_iter_dataset_records_round_trip_from_zip() -> None:
|
||||||
|
data = _zip_bytes({"revexp_66.xml": _SAMPLE_XML})
|
||||||
|
records = list(fol.iter_dataset_records(data, dataset="revexp", period=date(2026, 8, 25)))
|
||||||
|
assert len(records) == 3 # 1000 + (2000, 500)
|
||||||
|
inns = {r.inn for r in records}
|
||||||
|
assert inns == {"7707329152", "6660000001"}
|
||||||
|
assert all(r.dataset == "revexp" and r.period == date(2026, 8, 25) for r in records)
|
||||||
|
|
||||||
|
|
||||||
|
def test_iter_dataset_records_raises_when_zip_has_no_xml() -> None:
|
||||||
|
data = _zip_bytes({"readme.txt": "no xml here"})
|
||||||
|
with pytest.raises(ValueError, match="\\.xml"):
|
||||||
|
list(fol.iter_dataset_records(data, dataset="revexp", period=date(2026, 1, 1)))
|
||||||
|
|
||||||
|
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
# load_dataset — HTTP замокан (никакой живой сети)
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
def _make_client(zip_bytes: bytes, *, file_url: str, catalog_html: str) -> httpx.Client:
|
||||||
|
def handler(request: httpx.Request) -> httpx.Response:
|
||||||
|
if str(request.url) == file_url:
|
||||||
|
return httpx.Response(200, content=zip_bytes)
|
||||||
|
if "opendata/7707329152-" in str(request.url):
|
||||||
|
return httpx.Response(200, text=catalog_html)
|
||||||
|
return httpx.Response(404)
|
||||||
|
|
||||||
|
return httpx.Client(transport=httpx.MockTransport(handler))
|
||||||
|
|
||||||
|
|
||||||
|
def test_load_dataset_dry_run_does_not_write(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||||
|
file_url = "https://file.nalog.ru/opendata/7707329152-revexp/data-20260825-structure-1.zip"
|
||||||
|
catalog_html = f'<a href="{file_url}">zip</a>'
|
||||||
|
client = _make_client(b"unused", file_url=file_url, catalog_html=catalog_html)
|
||||||
|
|
||||||
|
db = MagicMock()
|
||||||
|
db.execute.return_value.first.return_value = None # get_loaded_version → None
|
||||||
|
|
||||||
|
result = fol.load_dataset(db, "revexp", client=client, dry_run=True)
|
||||||
|
assert result == {"dataset": "revexp", "skipped": 0, "records": 0, "upserted": 0}
|
||||||
|
# dry_run: единственный execute — SELECT версии (get_loaded_version), UPSERT не звался.
|
||||||
|
assert db.execute.call_count == 1
|
||||||
|
|
||||||
|
|
||||||
|
def test_load_dataset_skips_when_version_unchanged() -> None:
|
||||||
|
file_url = "https://file.nalog.ru/opendata/7707329152-revexp/data-20260825-structure-1.zip"
|
||||||
|
catalog_html = f'<a href="{file_url}">zip</a>'
|
||||||
|
client = _make_client(b"should-not-be-fetched", file_url=file_url, catalog_html=catalog_html)
|
||||||
|
|
||||||
|
db = MagicMock()
|
||||||
|
db.execute.return_value.first.return_value = (file_url, date(2026, 8, 25))
|
||||||
|
|
||||||
|
result = fol.load_dataset(db, "revexp", client=client)
|
||||||
|
assert result == {"dataset": "revexp", "skipped": 1, "records": 0, "upserted": 0}
|
||||||
|
|
||||||
|
|
||||||
|
def test_load_dataset_unknown_slug_raises() -> None:
|
||||||
|
db = MagicMock()
|
||||||
|
with pytest.raises(ValueError, match="неизвестный slug"):
|
||||||
|
fol.load_dataset(db, "not-a-real-dataset")
|
||||||
|
|
||||||
|
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
# Статические asserts по SQL (без БД)
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
def test_upsert_facts_sql_uses_cast_never_double_colon() -> None:
|
||||||
|
sql = re.sub(r"\s+", " ", str(fol._UPSERT_FACTS_SQL.text))
|
||||||
|
assert not re.search(r":\w+::", sql)
|
||||||
|
assert "CAST(:inn AS text)" in sql
|
||||||
|
assert "CAST(:period AS date)" in sql
|
||||||
|
assert "CAST(:value AS numeric)" in sql
|
||||||
|
assert "ON CONFLICT (inn, dataset, series, period) DO UPDATE SET" in sql
|
||||||
|
assert "IS DISTINCT FROM" in sql
|
||||||
|
|
||||||
|
|
||||||
|
def test_upsert_version_sql_uses_cast_and_conflict_target() -> None:
|
||||||
|
sql = re.sub(r"\s+", " ", str(fol._UPSERT_VERSION_SQL.text))
|
||||||
|
assert not re.search(r":\w+::", sql)
|
||||||
|
assert "ON CONFLICT (dataset) DO UPDATE SET" in sql
|
||||||
|
|
||||||
|
|
||||||
|
def test_upsert_records_uses_savepoint_and_module_does_not_commit() -> None:
|
||||||
|
src = inspect.getsource(fol.upsert_records)
|
||||||
|
assert "with db.begin_nested():" in src
|
||||||
|
module_src = inspect.getsource(fol)
|
||||||
|
assert "db.commit()" not in module_src
|
||||||
|
|
||||||
|
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
# fns_lookup.get_facts_by_inn — MagicMock db
|
||||||
|
# ─────────────────────────────────────────────────────────────────────────────
|
||||||
|
def test_get_facts_by_inn_groups_by_dataset() -> None:
|
||||||
|
db = MagicMock()
|
||||||
|
org = "ООО РОМАШКА"
|
||||||
|
p2026 = date(2026, 8, 25)
|
||||||
|
|
||||||
|
def row(dataset: str, series: str, period: date, value: int) -> dict[str, object]:
|
||||||
|
return {
|
||||||
|
"dataset": dataset,
|
||||||
|
"series": series,
|
||||||
|
"period": period,
|
||||||
|
"value": value,
|
||||||
|
"org_name": org,
|
||||||
|
}
|
||||||
|
|
||||||
|
db.execute.return_value.mappings.return_value.all.return_value = [
|
||||||
|
row("revexp", "СумДоход", p2026, 1000),
|
||||||
|
row("revexp", "СумРасход", p2026, 500),
|
||||||
|
row("sshr2019", "ССЧ", date(2019, 1, 1), 42),
|
||||||
|
]
|
||||||
|
result = fl.get_facts_by_inn(db, "7707329152")
|
||||||
|
assert set(result) == {"revexp", "sshr2019"}
|
||||||
|
assert len(result["revexp"]) == 2
|
||||||
|
assert result["sshr2019"][0]["series"] == "ССЧ"
|
||||||
|
|
||||||
|
|
||||||
|
def test_get_facts_by_inn_empty_inn_returns_empty_without_query() -> None:
|
||||||
|
db = MagicMock()
|
||||||
|
assert fl.get_facts_by_inn(db, "") == {}
|
||||||
|
assert fl.get_facts_by_inn(db, " ") == {}
|
||||||
|
db.execute.assert_not_called()
|
||||||
|
|
||||||
|
|
||||||
|
def test_get_facts_by_inn_no_rows_returns_empty_dict() -> None:
|
||||||
|
db = MagicMock()
|
||||||
|
db.execute.return_value.mappings.return_value.all.return_value = []
|
||||||
|
assert fl.get_facts_by_inn(db, "0000000000") == {}
|
||||||
|
|
||||||
|
|
||||||
|
_REAL_SHAPE_XML = """<?xml version="1.0" encoding="UTF-8"?>
|
||||||
|
<Файл ИдФайл="VO_OTKRDAN_5_9965_9965_20260825_cf1c08a2" ВерсФорм="4.01" \
|
||||||
|
ТипИнф="ОТКРДАННЫЕ5" КолДок="2">
|
||||||
|
<ИдОтпр><ФИООтв Фамилия="_" Имя="_"/></ИдОтпр>
|
||||||
|
<Документ ИдДок="dea28452-05f3-42bf-a08f-7a648c593473" ДатаДок="25.08.2026" \
|
||||||
|
ДатаСост="31.12.2025">
|
||||||
|
<СведНП НаимОрг="ОБЩЕСТВО С ОГРАНИЧЕННОЙ ОТВЕТСТВЕННОСТЬЮ "СПК"" \
|
||||||
|
ИННЮЛ="4205406898"/>
|
||||||
|
<СведДохРасх СумДоход="341864000.00" СумРасход="282224000.00"/>
|
||||||
|
</Документ>
|
||||||
|
<Документ ИдДок="28be840b-5db9-422b-b80c-3b678b808c96" ДатаДок="25.08.2026" \
|
||||||
|
ДатаСост="31.12.2025">
|
||||||
|
<СведНП НаимОрг="ООО ВАСИЛЁК" ИННЮЛ="6660000001"/>
|
||||||
|
<СведДохРасх СумДоход="1000.00" СумРасход="500.00"/>
|
||||||
|
</Документ>
|
||||||
|
</Файл>
|
||||||
|
"""
|
||||||
|
|
||||||
|
|
||||||
|
def test_iter_dataset_records_parses_real_fns_document_shape() -> None:
|
||||||
|
"""Боевая форма ФНС: ИНН на `СведНП`, показатели на СОСЕДНЕМ `СведДохРасх`.
|
||||||
|
|
||||||
|
Регресс на реальный дефект: харвестер требовал ИНН и числа на ОДНОМ элементе,
|
||||||
|
поэтому на настоящем дампе (revexp, data-20260825, 2118 XML) извлекал НОЛЬ
|
||||||
|
записей, а тесты на выдуманной фикстуре этого не показывали. Плюс безусловный
|
||||||
|
`elem.clear()` вычищал детей до закрытия `<Документ>`.
|
||||||
|
"""
|
||||||
|
data = _zip_bytes({"revexp_1.xml": _REAL_SHAPE_XML})
|
||||||
|
|
||||||
|
records = list(fol.iter_dataset_records(data, dataset="revexp", period=date(2026, 8, 25)))
|
||||||
|
|
||||||
|
assert len(records) == 4 # 2 организации × (СумДоход, СумРасход)
|
||||||
|
assert {r.inn for r in records} == {"4205406898", "6660000001"}
|
||||||
|
assert {r.series for r in records} == {"СумДоход", "СумРасход"}
|
||||||
|
# период берётся из ДатаСост документа, а не из переданного fallback'а
|
||||||
|
assert {r.period for r in records} == {date(2025, 12, 31)}
|
||||||
|
first = next(r for r in records if r.inn == "4205406898" and r.series == "СумДоход")
|
||||||
|
assert first.value == 341864000.0
|
||||||
|
assert first.org_name is not None and "СПК" in first.org_name
|
||||||
|
|
||||||
|
|
||||||
|
def test_service_ids_are_not_harvested_as_facts() -> None:
|
||||||
|
"""ВерсФорм/КолДок/ДатаДок — служебные, фактами быть не должны."""
|
||||||
|
data = _zip_bytes({"revexp_1.xml": _REAL_SHAPE_XML})
|
||||||
|
series = {
|
||||||
|
r.series for r in fol.iter_dataset_records(data, dataset="revexp", period=date(2026, 8, 25))
|
||||||
|
}
|
||||||
|
assert not series & {"ВерсФорм", "КолДок", "ДатаДок", "ДатаСост", "ИННЮЛ"}
|
||||||
Loading…
Add table
Reference in a new issue