gendesign/tradein-mvp/backend/app/services/fns_opendata_loader.py
lekss361 e0bef636e6
All checks were successful
Deploy Trade-In / changes (push) Successful in 12s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 4m17s
Deploy Trade-In / build-backend (push) Successful in 1m3s
Deploy Trade-In / deploy (push) Successful in 7m11s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / perimeter-smoke (push) Successful in 11s
feat(tradein): bulk-дампы открытых данных ФНС по юрлицам + lookup по ИНН (#3429)
2026-09-08 22:29:10 +00:00

497 lines
25 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.

"""ФНС 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
}