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")
Helpful? Dislike 0 Log in to react