Programming
2. init.py
from . import config, io_layer, enrich, marts, quality, i18n # noqa
__version__ = "0.1.0"
"""
Raw file -> partitioned parquet lake -> typed snapshot frames.
Design rule: a month is a partition. Re-ingesting a month overwrites that
partition and nothing else, so "the September file was wrong, here it is again"
is a one-line fix instead of a rebuild.
"""
from __future__ import annotations
import shutil
from pathlib import Path
import numpy as np
import pandas as pd
from . import config as C
CANONICAL_COLUMNS = [
"contragent", "anketa", "filial", "tobo", "valuta", "dogovor_sana",
"vidacha_sana", "department", "department_number", "passport",
"contragent_name", "brutto_95_summa", "brutto_summa",
"protsent_16309_summa", "protsent_16377_summa", "protsent_16379_summa",
"spisat_95413_summa", "spisat_91501_summa", "reserve_summa",
"vidacha_tekushiy_summa", "pogashen_tekushiy_summa", "maks_dni",
"kategoriya_origin", "kategoriya_npl_90", "straxovka_summa",
"straxovka_brutto_95_summa", "kategoriya_straxovka", "sud", "avto",
"report_date", "all_time_given_amount",
]
_ALIASES = {
"vidacha_tm_summa": "vidacha_tekushiy_summa",
"vidacha_alltime_summa": "all_time_given_amount",
"pogashen_tm_summa": "pogashen_tekushiy_summa",
}
def _read_any(path: Path, **kw) -> pd.DataFrame:
s = path.suffix.lower()
if s == ".parquet":
return pd.read_parquet(path, **kw)
if s in (".csv", ".txt"):
return pd.read_csv(path, low_memory=False, **kw)
if s == ".xlsb":
return pd.read_excel(path, engine="pyxlsb", **kw)
if s in (".xlsx", ".xlsm"):
return pd.read_excel(path, **kw)
raise ValueError(f"unsupported source format: {path.name}")
def normalise(df: pd.DataFrame) -> pd.DataFrame:
"""Column names, dtypes and the handful of encoding quirks in the source."""
df = df.rename(columns=lambda c: str(c).strip().lower().replace(" ", "_"))
df = df.rename(columns=_ALIASES)
missing = [c for c in CANONICAL_COLUMNS if c not in df.columns]
if missing:
raise KeyError(f"source is missing required columns: {missing}")
df = df[list(CANONICAL_COLUMNS)].copy()
for c in C.DATE_COLS:
df[c] = pd.to_datetime(df[c], errors="coerce")
for c in C.AMOUNT_COLS:
df[c] = pd.to_numeric(df[c], errors="coerce").fillna(0.0).astype("float64")
df["maks_dni"] = pd.to_numeric(df["maks_dni"], errors="coerce").fillna(0).astype("int32")
for c in ("kategoriya_origin", "kategoriya_npl_90", "kategoriya_straxovka"):
df[c] = pd.to_numeric(df[c], errors="coerce").fillna(0).astype("int8")
df["valuta"] = pd.to_numeric(df["valuta"], errors="coerce").fillna(0).astype("int32")
for c in ("anketa", "contragent"):
df[c] = pd.to_numeric(df[c], errors="coerce").astype("Int64")
for c in ("filial", "tobo", "department", "department_number", "passport",
"contragent_name", "sud", "avto"):
df[c] = df[c].astype("string").str.strip()
# `avto` mixes the integer 0 with two labels — make it one clean vocabulary.
df["avto"] = (df["avto"].fillna("0")
.replace({"0": "non_auto", "0.0": "non_auto", "": "non_auto",
"nan": "non_auto", "None": "non_auto"}))
df["sud"] = df["sud"].fillna("norm")
return df
def ingest(source: Path | str, *, overwrite: bool = True,
verbose=None) -> list[pd.Timestamp]:
"""
Load a raw file (full history or a single new month) into the lake.
Streams one report_date at a time rather than reading the whole file. A
full 8M-row history is 3–6 GB in pandas if loaded at once; one month is
~250k rows, which is comfortable anywhere. Returns the report_dates written.
"""
source = Path(source)
log = verbose or (lambda *_: None)
dates = _scan_report_dates(source)
if not dates:
raise ValueError(f"no usable report_date values in {source.name}")
written: list[pd.Timestamp] = []
for rd in dates:
target = C.LAKE_DIR / f"report_date={rd:%Y-%m-%d}"
if target.exists() and not overwrite:
log(f" skip {rd:%Y-%m} (already in lake)")
continue
part = normalise(_read_one_date(source, rd))
if part.empty:
continue
if target.exists():
shutil.rmtree(target)
target.mkdir(parents=True)
part.drop(columns=["report_date"]).to_parquet(
target / "part.parquet", index=False, compression="zstd")
written.append(rd)
log(f" {rd:%Y-%m}: {len(part):,} rows")
del part
return written
def _scan_report_dates(path: Path) -> list[pd.Timestamp]:
"""Read only the report_date column, so a huge file costs almost nothing."""
s = path.suffix.lower()
if s == ".parquet":
col = _report_date_column_name(path)
vals = pd.read_parquet(path, columns=[col])[col]
else:
df = _read_any(path)
df = df.rename(columns=lambda c: str(c).strip().lower().replace(" ", "_"))
vals = df["report_date"]
vals = pd.to_datetime(vals, errors="coerce").dropna().unique()
return sorted(pd.Timestamp(v).normalize() for v in vals)
def _report_date_column_name(path: Path) -> str:
import pyarrow.parquet as pq
names = pq.ParquetFile(path).schema_arrow.names
for n in names:
if str(n).strip().lower().replace(" ", "_") == "report_date":
return n
raise KeyError(f"{path.name} has no report_date column")
def _read_one_date(path: Path, rd: pd.Timestamp) -> pd.DataFrame:
"""One month's rows. Parquet gets predicate pushdown; other formats don't."""
if path.suffix.lower() == ".parquet":
col = _report_date_column_name(path)
return pd.read_parquet(path, filters=[(col, "==", rd.to_pydatetime())])
df = _read_any(path)
df = df.rename(columns=lambda c: str(c).strip().lower().replace(" ", "_"))
return df[pd.to_datetime(df["report_date"], errors="coerce").eq(rd)]
def available_snapshots() -> list[pd.Timestamp]:
dates = [pd.Timestamp(p.name.split("=", 1)[1]) for p in C.LAKE_DIR.glob("report_date=*")]
return sorted(dates)
def load_snapshot(report_date, columns: list[str] | None = None) -> pd.DataFrame:
rd = pd.Timestamp(report_date)
path = C.LAKE_DIR / f"report_date={rd:%Y-%m-%d}" / "part.parquet"
if not path.exists():
raise FileNotFoundError(f"no snapshot for {rd:%Y-%m-%d}")
cols = [c for c in columns if c != "report_date"] if columns else None
df = pd.read_parquet(path, columns=cols)
df["report_date"] = rd
return df
def iter_snapshots(columns: list[str] | None = None, dates=None):
for rd in (dates or available_snapshots()):
yield rd, load_snapshot(rd, columns)
# --------------------------------------------------------------------------
# FX
# --------------------------------------------------------------------------
def load_fx() -> pd.DataFrame:
"""
fx_rates.csv: report_date,valuta,rate (rate = UZS per 1 unit).
Missing months fall back to config.FX_FALLBACK so the pipeline never
silently drops FX exposure.
"""
snaps = available_snapshots()
rows = []
if C.FX_PATH.exists():
fx = pd.read_csv(C.FX_PATH)
fx["report_date"] = pd.to_datetime(fx["report_date"])
rows.append(fx[["report_date", "valuta", "rate"]])
grid = pd.DataFrame(
[(d, c) for d in snaps for c in C.FX_FALLBACK],
columns=["report_date", "valuta"])
grid["rate"] = grid["valuta"].map(C.FX_FALLBACK)
rows.append(grid)
fx = pd.concat(rows, ignore_index=True)
fx = fx.drop_duplicates(subset=["report_date", "valuta"], keep="first")
fx["valuta"] = fx["valuta"].astype("int32")
return fx
def fx_factor(df: pd.DataFrame, fx: pd.DataFrame, mode: str = "nominal") -> np.ndarray:
if mode == "constant":
anchor = fx[fx["report_date"] == pd.Timestamp(C.FX_CONSTANT_DATE)]
m = dict(zip(anchor["valuta"], anchor["rate"])) or C.FX_FALLBACK
return df["valuta"].map(m).fillna(1.0).to_numpy(dtype="float64")
key = list(zip(df["report_date"], df["valuta"]))
m = {(r, c): v for r, c, v in zip(fx["report_date"], fx["valuta"], fx["rate"])}
return np.array([m.get(k, C.FX_FALLBACK.get(k[1], 1.0)) for k in key], dtype="float64")
PO
powerty
Author
· Staff
Aug. 26, 2026
Aug. 26, 2026
3
Views
0
Likes
4m
Read