"""ФНС opendata loader: доходы/расходы, ССЧ, недоимка, спецрежимы юрлиц → lookup по ИНН. CONTEXT: открытых данных ЕГРН по правообладателям-физлицам не существует (218-ФЗ ст. 62) — это жёсткий блокер. По ЮРЛИЦАМ данные открыты, лицензия ФНС (nalog.gov.ru/opendata) разрешает переработку и перераспространение — легальный обход для того среза, где он в принципе доступен. Наборы (маска slug'а: 7707329152-): revexp — доходы и расходы; sshr2019 — среднесписочная численность; debtam — недоимка и задолженность; snr — спецрежимы. ЕГРЮЛ-данных здесь НЕТ: ни адреса, ни ОКВЭД, ни учредителей — только ИНН + наименование + показатели набора. ПОТРЕБИТЕЛЯ У ЭТОГО СЕРВИСА ПОКА НЕТ. Этот модуль (+ fns_lookup.py) — только загрузка и lookup по ИНН. Подключение к продукту (например, скоринг застройщика/УК в оценке) сюда сознательно не входит и не реализовано — это решается отдельной задачей вне текущего PR. `estimator.py` не тронут. ИСТОЧНИК: https://www.nalog.gov.ru/opendata/7707329152-/ — 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 }