Загрузка данных


"""Загрузка из DWH, кластеризация, запись в core (DAG: delete_full_insert)."""
from __future__ import annotations

from datetime import datetime
from typing import Any

import pandas as pd
from processing_class.mssql_pyodbc_hook import MsSqlPyodbcHook

from service import decorator_func as dec
from sql_script import custom_select_source as custom_select

from processing_class.nomenclature_fuzzy_groups.cluster import cluster_fuzzy, fuzzy_block
from processing_class.nomenclature_fuzzy_groups.constants import (
    FUZZY_BLOCK_PREFIX_LEN,
    FUZZY_CORE_THRESHOLD,
    SINGLE_TOKEN_TAIL_MIN_RATIO,
    VARIANT_TAIL_SPLIT_THRESHOLD,
)
from processing_class.nomenclature_fuzzy_groups.naming import normalize_nomenclature, parse_name


def load_nomenclature(hook: Any) -> pd.DataFrame:
    return hook.get_pandas_df(sql=custom_select.nomenclature_fuzzy_groups)


def nomen_group_cluster_max(s: pd.Series) -> object:
    strs = [str(x).strip() for x in s.dropna() if str(x).strip()]
    if not strs:
        return pd.NA
    num = pd.to_numeric(pd.Series(strs, dtype=object), errors="coerce")
    if num.notna().all():
        mx = num.max()
        tied = [t for t, n in zip(strs, num) if n == mx]
        return max(tied)
    return max(strs)


def assign_cluster_nomen_group(df: pd.DataFrame) -> pd.DataFrame:
    out = df.copy()
    out["nomen_group_cluster"] = out.groupby("fuzzy_cluster_id", sort=False)["nomen_group"].transform(
        nomen_group_cluster_max
    )
    return out


def process_dataframe(df: pd.DataFrame) -> pd.DataFrame:
    """Кластеризация по загруженному DataFrame (без SQL)."""
    out = df.copy()
    out["normalized_name"] = out["nomenclature_name"].map(
        lambda x: normalize_nomenclature(x if pd.notna(x) else "")
    )
    parsed = [parse_name(n) for n in out["normalized_name"].tolist()]
    out["name_core"] = [p.core for p in parsed]
    packs = [p.pack for p in parsed]
    vars_only = [p.variant for p in parsed]
    out["tail_pack"] = packs
    out["variant_tail"] = [f"{v} {p}".strip() if p else v for v, p in zip(vars_only, packs)]
    out["fuzzy_block"] = [fuzzy_block(n, FUZZY_BLOCK_PREFIX_LEN) for n in out["normalized_name"]]
    out["fuzzy_cluster_id"] = cluster_fuzzy(
        out["name_core"].tolist(),
        vars_only,
        packs,
        out["fuzzy_block"].tolist(),
        FUZZY_CORE_THRESHOLD,
        VARIANT_TAIL_SPLIT_THRESHOLD,
        SINGLE_TOKEN_TAIL_MIN_RATIO,
    )
    out = assign_cluster_nomen_group(out)
    return out.sort_values(
        by=["nomen_group_cluster", "nomenclature_code_first"],
        ascending=[True, True],
        na_position="last",
        ignore_index=True,
    )


class NomenclatureFuzzyGroups:
    """Расчёт fuzzy-групп и полная перезапись таблицы в core (TRUNCATE + INSERT)."""

    def __init__(
        self,
        connect_id: str,
        db_core: str,
        schema_core: str,
        table_name_core: str,
        fields_core: list,
    ):
        self._mssql_hook = MsSqlPyodbcHook(mssql_conn_id=connect_id)
        self._fields_core = fields_core
        self._table_name_core_path = f"[{db_core}].[{schema_core}].[{table_name_core}]"

    def _bulk_insert_core(self, df: pd.DataFrame, commit_every: int = 5000) -> None:
        cols = self._fields_core
        col_list = ", ".join(f"[{c}]" for c in cols) + ", [dlm$]"
        placeholders = ", ".join("?" * len(cols)) + ", ?"
        insert_sql = (
            f"INSERT INTO {self._table_name_core_path} ({col_list}) "
            f"VALUES ({placeholders})"
        )
        dlm = datetime.now()
        rows = [
            tuple(None if pd.isna(v) else v for v in row) + (dlm,)
            for row in df[cols].itertuples(index=False, name=None)
        ]
        if not rows:
            return

        conn = self._mssql_hook.get_conn()
        try:
            cur = conn.cursor()
            cur.fast_executemany = True
            for start in range(0, len(rows), commit_every):
                cur.executemany(insert_sql, rows[start:start + commit_every])
            conn.commit()
        finally:
            conn.close()

    @dec.task_python_operator
    @dec.timer('Nomenclature fuzzy groups delete_full_insert')
    def delete_full_insert(self, **kwargs):
        result = process_dataframe(load_nomenclature(self._mssql_hook))[self._fields_core]
        self._mssql_hook.get_records(f"TRUNCATE TABLE {self._table_name_core_path}")
        self._bulk_insert_core(result)
        print(f'Записано в {self._table_name_core_path}: {len(result):,} строк')
--------------------------------------------------------------------------------------------------------------------------------------------------------


from airflow.models.dag import DAG

from processing_class.nomenclature_fuzzy_groups.processing_nomenclature import NomenclatureFuzzyGroups
from service.notify_failure_cooldown import notify_failure_cooldown
import pendulum

default_args = {
    'owner': 'nomenclature',
    'depends_on_past': False,
    'retries': 0,
    'on_failure_callback': notify_failure_cooldown(60),
    'start_date': pendulum.datetime(2026, 4, 1, tz='Europe/Moscow'),
}

connector_id = 'mssql_dwh'
db_core = 'dwh'
schema_core = 'dbo'

with DAG(
    dag_id='nomenclature_fuzzy_groups',
    default_args=default_args,
    tags=['nomenclature', 'fuzzy', 'dwh'],
    schedule_interval='0 20 * * *',
    catchup=False
):
    table_name_core = 'nomenclature_fuzzy_groups'
    fields_core = [
        'nomenclature_name', 'normalized_name', 'nomen_group', 'nomen_group_cluster',
        'group_buyer_cnt', 'group_rebrand_cnt', 'nomenclature_code_customer_first',
        'rebrand_sku_erp_first', 'nomenclature_code_first',
    ]

    main = NomenclatureFuzzyGroups(
        connector_id, db_core, schema_core, table_name_core, fields_core,
    )

    main.delete_full_insert()